myc

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

commit b72c9fe073c3776294030003eb62b743c660927f
parent 89285d7a1f3c045fc952ac2caadd303dd646254b
Author: triesap <tyson@radroots.org>
Date:   Sun, 23 Aug 2026 16:10:36 +0000

runtime: own the Myc daemon graph

- wire the fixed admin, operations, relay, provider, and delivery roles
- preserve exact request replay and bounded full-jitter delivery retries
- publish real status and coordinate one-deadline graceful shutdown
- freeze configuration, source-lock, package, and API evidence

Diffstat:
MCargo.lock | 12++++++++++++
MCargo.toml | 3++-
MREADME | 43+++++++++++++++++++++++++++----------------
Mcontracts/api_baselines/myc.txt | 11+++++++++++
Mcontracts/services_hardening/config.v1.example.toml | 5+++++
Mcontracts/services_hardening/config.v1.schema.json | 29++++++++++++++++++++++++++++-
Mradroots.service.source-lock.v2.toml | 2+-
Msrc/delivery_worker.rs | 111+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++--------------------
Msrc/lib.rs | 8+++++++-
Msrc/main.rs | 53++++++++++++++++++++++++++++++++++++++++++++++++++++-
Msrc/nip46_replay.rs | 5+++++
Msrc/nip46_wave_080_a.rs | 2+-
Msrc/nip46_work.rs | 12++++++++++++
Msrc/process_v1.rs | 115+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++------------
Msrc/process_v1_unsupported.rs | 13+++++++++++++
Msrc/provider_contract.rs | 56++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Asrc/runtime_graph.rs | 2761+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Asrc/runtime_nip46.rs | 718+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Asrc/runtime_signal.rs | 50++++++++++++++++++++++++++++++++++++++++++++++++++
Msrc/runtime_supervision.rs | 4++++
Msrc/state_admin.rs | 250++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-
Msrc/state_connection.rs | 414++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-
Msrc/state_delivery.rs | 104+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Msrc/state_discovery.rs | 85+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++--
Msrc/state_governance.rs | 207+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++--
Msrc/state_request.rs | 10+++++-----
Msrc/state_response.rs | 22++++++++++++++++++++++
Msrc/transport_nostr_adapter.rs | 226++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++---
Mtests/build_policy.rs | 9++++++++-
Mtests/package_boundary.rs | 49+++++++++++++++++++++++++++++++++++++++++++++++--
Mtests/services_hardening_config_contract.rs | 27+++++++++++++++++++++++++++
Mtests/services_hardening_diagnostics.rs | 10+++++-----
Mtests/services_hardening_legacy_removal.rs | 5++---
Mtests/services_hardening_native_release.rs | 2+-
Mtests/services_hardening_process.rs | 18++++++++++++++++++
Mtests/services_hardening_state_metadata.rs | 16+++++++++++++++-
36 files changed, 5364 insertions(+), 103 deletions(-)

diff --git a/Cargo.lock b/Cargo.lock @@ -1343,6 +1343,7 @@ dependencies = [ "jsonschema", "nostr", "radroots_event_codec", + "radroots_identity", "radroots_nostr", "radroots_nostr_connect", "radroots_runtime_paths", @@ -2266,6 +2267,16 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "0fda2ff0d084019ba4d7c6f371c95d8fd75ce3524c3cb8fb653a3023f6323e64" [[package]] +name = "signal-hook-registry" +version = "1.4.8" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "c4db69cba1110affc0e9f7bcd48bbf87b3f4fc7c61fc9155afd4c469eb3d6c1b" +dependencies = [ + "errno", + "libc", +] + +[[package]] name = "slab" version = "0.4.12" source = "registry+https://github.com/rust-lang/crates.io-index" @@ -2560,6 +2571,7 @@ dependencies = [ "libc", "mio", "pin-project-lite", + "signal-hook-registry", "socket2", "tokio-macros", "windows-sys 0.61.2", diff --git a/Cargo.toml b/Cargo.toml @@ -56,6 +56,7 @@ futures-executor = "0.3" hex = "0.4" jsonschema = { version = "0.48.1", default-features = false } nostr = { version = "0.44.2", features = ["nip04", "nip44", "nip46", "nip49"] } +radroots_identity = { git = "https://github.com/radrootslabs/lib", rev = "7d7b454b4c9ed86569671993bd03ca868b676665", version = "=0.1.0-alpha" } radroots_nostr = { git = "https://github.com/radrootslabs/lib", rev = "7d7b454b4c9ed86569671993bd03ca868b676665", version = "=0.1.0-alpha", features = ["events"] } radroots_nostr_connect = { git = "https://github.com/radrootslabs/lib", rev = "7d7b454b4c9ed86569671993bd03ca868b676665", version = "=0.1.0-alpha" } radroots_event_codec = { git = "https://github.com/radrootslabs/lib", rev = "7d7b454b4c9ed86569671993bd03ca868b676665", version = "=0.1.0-alpha", features = ["json"] } @@ -71,7 +72,7 @@ serde_json = { version = "1.0", features = ["raw_value"] } sha2 = "0.10" sqlx = { version = "0.9.0", default-features = false, features = ["derive", "sqlite-bundled"] } rustix = { version = "1", features = ["fs", "process", "std"] } -tokio = { version = "1.48", default-features = false, features = ["io-util", "macros", "net", "rt-multi-thread", "sync", "time"] } +tokio = { version = "1.48", default-features = false, features = ["io-util", "macros", "net", "rt-multi-thread", "signal", "sync", "time"] } toml = "0.8" url = "2.5" zeroize = "1.8" diff --git a/README b/README @@ -64,11 +64,10 @@ init|validate|show|schema|apply`; `state init|status|backup|restore|verify|migra is an offline create-new/config-apply lifecycle and is intentionally absent from the live CLI and Unix-admin inventories. The process binary uses only this parser and consumes its sealed execution plan -exactly once. Every non-daemon command now reaches its governed config, state, -identity, Unix-admin, or doctor authority. `run` remains explicitly unavailable -until the next ordered unit installs the authoritative daemon graph; no admitted -but unimplemented command can return success and no prototype command is used -as a fallback. +exactly once. Every command reaches its governed config, state, identity, +Unix-admin, doctor, or daemon authority. `run` installs the authoritative +bounded daemon graph; no admitted command can return success through a +prototype or unavailable-operation fallback. The selected absolute config path is opened no-follow through its retained parent descriptor. The loader requires a regular, single-link, @@ -149,14 +148,16 @@ codes. It accepts no caller text, path, SQL, raw cause, identity, relay URL, credential, secret, or decrypted content. Every public Myc error remains a crate-owned source-free classification, so ordinary `Display`, `Debug`, and whole-chain traversal cannot bypass redaction. Current binary dispatch reports -invalid input as exit 2 and the intentionally unavailable daemon graph as exit -3. The Myc runtime now owns one sealed, bounded critical-task graph over -the shared supervisor: task error, panic, join failure, unexpected cancellation, -or success before cancellation cancels peers, joins every task, and maps to the -fixed unexpected-internal process result. Tasks receive only a cooperative -cancellation observer; task names and handles remain internal. Step 159 owns -the process panic hook, signal installation, externally requested and forced -shutdown, the configured grace deadline, and durability-aware ordered drain. +invalid input as exit 2 and service dependency unavailability as exit 3. The +Myc runtime owns one sealed, bounded critical-task graph over the shared supervisor: +task error, panic, join failure, unexpected cancellation, or critical success +before cancellation coordinates shutdown, joins every task, and maps to the +fixed nonzero process result. Tasks receive only a cooperative cancellation +observer; task names and handles remain internal. The binary owns signal +installation. The graph uses one configured absolute graceful-shutdown deadline +for mutation rejection, ingress and operations drain, +recoverable-work persistence, network and SQLite close, and socket close; it +adds no hidden second cleanup budget. `parse_myc_config_v1` caps original bytes before decoding, checks the schema header before closed contract admission, rejects duplicate, null, unknown, and @@ -229,9 +230,19 @@ persists Submitted immediately before execution. Accepted, rejected, transport-failed, and unknown acknowledgements remain distinct; cancellation or lost acknowledgement after Submitted is durably unknown, and retries never alter the committed bytes. Provider or relay work never occurs inside a SQLite -transaction. Unit 14's secure command executor uses these exact provider and -read-only relay probe boundaries; runtime task-graph wiring and startup -handshakes remain the next ordered Step 159 unit. +transaction. The secure command executor and production daemon graph use these +same provider, relay, and exact-byte delivery boundaries. + +The production `run` path owns the exact five-role bounded graph: +`admin_server`, optional `operations_server`, `relay_ingress`, +`provider_dispatch`, and `delivery_outbox`. It creates no task per request or +relay. Required relay subscriptions and provider handshakes complete before +Ready; bounded reconnect publishes Unready when a required dependency is lost, +while optional operations loss publishes Degraded. NIP-46 dispatch verifies, +decrypts, admits, authorizes, executes providers outside transactions, +re-encrypts in the verified request context, signs through the user provider, +and atomically commits completion, exact response bytes, immutable targets, +and initial outbox state. Completed replay never re-executes a provider. Schema v8 adds immutable NIP-46 operation-completion evidence. The Step 147 integration checkpoint binds each durable request to its stable operation and diff --git a/contracts/api_baselines/myc.txt b/contracts/api_baselines/myc.txt @@ -578,6 +578,13 @@ impl myc::MycProcessResult pub const fn myc::MycProcessResult::code(self) -> &'static str pub fn myc::MycProcessResult::exit_code(self) -> std::process::ExitCode pub const fn myc::MycProcessResult::exit_code_u8(self) -> u8 +pub enum myc::MycProcessSignal +pub myc::MycProcessSignal::Interrupt +pub myc::MycProcessSignal::Terminate +impl myc::MycProcessSignal +pub const fn myc::MycProcessSignal::as_str(self) -> &'static str +impl core::fmt::Display for myc::MycProcessSignal +pub fn myc::MycProcessSignal::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result pub enum myc::MycProviderCapability pub myc::MycProviderCapability::Describe pub myc::MycProviderCapability::Nip04Decrypt @@ -2266,11 +2273,14 @@ pub trait myc::MycAdminHandler: core::marker::Send + core::marker::Sync + 'stati pub fn myc::MycAdminHandler::handle<'a>(&'a self, myc::MycAdminRequestDocument) -> myc::MycAdminFuture<'a> pub trait myc::MycDoctorProbe: core::marker::Send + core::marker::Sync pub fn myc::MycDoctorProbe::probe(&self, myc::MycDoctorCheckDefinition) -> myc::MycDoctorFuture<'_> +pub trait myc::MycProcessSignalSource: core::marker::Send +pub fn myc::MycProcessSignalSource::next_signal(&mut self) -> myc::MycProcessSignalFuture<'_> pub fn myc::admit_myc_nip46_event(myc::MycNip46AdmissionLimits, &[u8]) -> core::result::Result<myc::MycBoundedNip46Event, myc::MycNip46AdmissionError> pub fn myc::admit_myc_nip46_request(myc::MycNip46AdmissionLimits, &[u8]) -> core::result::Result<myc::MycBoundedNip46Request, myc::MycNip46AdmissionError> pub fn myc::bind_myc_nip46_replay(myc::MycVerifiedNip46Event, myc::MycVerifiedNip46Request) -> core::result::Result<myc::MycReplayBoundNip46Request, myc::MycSignerRequestError> pub fn myc::build_myc_admin_router<H>(alloc::sync::Arc<H>) -> core::result::Result<myc::MycAdminRouter, myc::MycAdminRouterError> where H: myc::MycAdminHandler pub fn myc::execute_myc_cli_v1(myc::MycCliInvocationV1) -> myc::MycProcessResult +pub fn myc::execute_myc_cli_v1_with_signal_source<F, S>(myc::MycCliInvocationV1, F) -> myc::MycProcessResult where F: core::ops::function::FnOnce() -> core::option::Option<S>, S: myc::MycProcessSignalSource + 'static pub async fn myc::finalize_myc_state_restore(myc::MycStagedStateRestore) -> core::result::Result<(), myc::MycStateMaintenanceError> pub fn myc::initialize_myc_config_document(&myc::MycRuntimeContext, &[u8]) -> core::result::Result<myc::MycConfigDocumentV1, myc::MycConfigLoadError> pub async fn myc::initialize_myc_state(&myc::MycRuntimeContext, &myc::MycStateMetadata, radroots_service_sqlite::migration::MigrationAppliedAtUnixSeconds, &radroots_service_sqlite::migration::MigrationBuildIdentity) -> core::result::Result<(), myc::MycStateHostError> @@ -2303,3 +2313,4 @@ pub fn myc::verify_myc_nip46_request(myc::MycBoundedNip46Request) -> core::resul pub fn myc::verify_myc_state_backup(&[u8], radroots_service_sqlite::backup::manifest::BackupManifestSha256, &std::path::Path, &myc::MycStateMetadata, core::num::nonzero::NonZeroU64) -> core::result::Result<myc::MycVerifiedStateBackup, myc::MycStateMaintenanceError> pub type myc::MycAdminFuture<'a> = core::pin::Pin<alloc::boxed::Box<(dyn core::future::future::Future<Output = core::result::Result<myc::MycAdminResponseDocument, myc::MycAdminHandlerError>> + core::marker::Send + 'a)>> pub type myc::MycDoctorFuture<'a> = core::pin::Pin<alloc::boxed::Box<(dyn core::future::future::Future<Output = myc::MycDoctorObservation> + core::marker::Send + 'a)>> +pub type myc::MycProcessSignalFuture<'a> = core::pin::Pin<alloc::boxed::Box<(dyn core::future::future::Future<Output = core::option::Option<myc::MycProcessSignal>> + core::marker::Send + 'a)>> diff --git a/contracts/services_hardening/config.v1.example.toml b/contracts/services_hardening/config.v1.example.toml @@ -59,6 +59,11 @@ authentication = "required" [transport] connect_deadline_ms = 10000 +[transport.ingress] +subscription_deadline_ms = 30000 +maximum_past_seconds = 120 +maximum_future_seconds = 120 + [transport.delivery_policy] mode = "all_required" diff --git a/contracts/services_hardening/config.v1.schema.json b/contracts/services_hardening/config.v1.schema.json @@ -272,13 +272,40 @@ "transport": { "type": "object", "additionalProperties": false, - "required": ["delivery_policy"], + "required": ["delivery_policy", "ingress"], "properties": { "connect_deadline_ms": { "type": "integer", "minimum": 1, "maximum": 30000, "default": 10000, "x-radroots-default-source": "engineering_safety" }, "delivery_policy": { "$ref": "#/$defs/delivery_policy" }, + "ingress": { "$ref": "#/$defs/ingress" }, "publish_retry": { "$ref": "#/$defs/publish_retry" } } }, + "ingress": { + "type": "object", + "additionalProperties": false, + "required": [ + "subscription_deadline_ms", + "maximum_past_seconds", + "maximum_future_seconds" + ], + "properties": { + "subscription_deadline_ms": { + "type": "integer", + "minimum": 1000, + "maximum": 300000 + }, + "maximum_past_seconds": { + "type": "integer", + "minimum": 0, + "maximum": 3600 + }, + "maximum_future_seconds": { + "type": "integer", + "minimum": 0, + "maximum": 3600 + } + } + }, "permission": { "type": "string", "minLength": 1, diff --git a/radroots.service.source-lock.v2.toml b/radroots.service.source-lock.v2.toml @@ -7,7 +7,7 @@ architecture = "radroots.crates.release.v2" workspace_catalog_sha256 = "deca0c080deae187ff8186c0708903e42f41ea57f77c5f91581e23aa561164a4" version = "0.1.0-alpha" source_archive_sha256 = "b425371c134be96cce46b37f7035d6212f1efe8cff50bef366631ba5632991b0" -cargo_lock_sha256 = "a49a45076e3fecb49cdf8d98ea1082f62d7025cbe5cd9359df1e7b89e36df6d4" +cargo_lock_sha256 = "8599fbc43a79dd13b0d496e86b7aefe2e6c863e016ac8a0d5612b7eb7cd979d6" rust_version = "1.97.1" host_feature_profile = "service-host" diff --git a/src/delivery_worker.rs b/src/delivery_worker.rs @@ -1,10 +1,5 @@ //! Durable delivery orchestration over exact committed event bytes. -#![allow( - dead_code, - reason = "Step 159 Unit 12 seals the worker before Unit 15 runtime graph wiring" -)] - use core::fmt; use std::error::Error; @@ -14,9 +9,10 @@ use crate::transport_nostr_adapter::{ MycNostrDeliveryAdapter, MycRelayAdapter, MycRelayExecutionOutcome, }; use crate::{ - MycConfigDocumentV1, MycDeliveryAttemptNonce, MycDeliveryAttemptOutcome, MycDeliveryClaim, - MycDeliveryJobId, MycDeliveryJobRecord, MycDeliveryRelayId, MycDeliveryRetryJitter, - MycDeliverySourceKind, MycDeliveryTimeUnixMs, MycStateRepository, MycTaskCancellation, + MycConfigDocumentV1, MycDeliveryAttemptNonce, MycDeliveryAttemptOutcome, + MycDeliveryAttemptRecord, MycDeliveryClaim, MycDeliveryJobId, MycDeliveryJobRecord, + MycDeliveryRelayId, MycDeliveryRetryJitter, MycDeliverySourceKind, MycDeliveryTimeUnixMs, + MycStateRepository, MycTaskCancellation, }; #[derive(Clone, Copy, Debug, PartialEq, Eq)] @@ -31,6 +27,7 @@ pub(crate) struct MycDeliveryWorkerError { } impl MycDeliveryWorkerError { + #[cfg(test)] pub(crate) const fn kind(&self) -> MycDeliveryWorkerErrorKind { self.kind } @@ -80,7 +77,7 @@ pub(crate) struct MycDeliveryExecutionEvidence { pub(crate) claimed_at: MycDeliveryTimeUnixMs, pub(crate) submitted_at: MycDeliveryTimeUnixMs, pub(crate) observed_at: MycDeliveryTimeUnixMs, - pub(crate) retry_jitter: MycDeliveryRetryJitter, + pub(crate) retry_entropy: [u8; 8], } impl fmt::Debug for MycDeliveryExecutionEvidence { @@ -164,10 +161,10 @@ async fn run_with_adapter<A: MycRelayAdapter>( Err(_) => { return persist_outcome( repository, - job_id, relay_id, - attempt.id(), MycDeliveryAttemptOutcome::TransportFailed, + &job, + &attempt, evidence, ) .await; @@ -176,10 +173,10 @@ async fn run_with_adapter<A: MycRelayAdapter>( if cancellation.is_cancelled() { return persist_outcome( repository, - job_id, relay_id, - attempt.id(), MycDeliveryAttemptOutcome::TransportFailed, + &job, + &attempt, evidence, ) .await; @@ -200,15 +197,7 @@ async fn run_with_adapter<A: MycRelayAdapter>( MycDeliveryAttemptOutcome::UnknownAcknowledgement } }; - persist_outcome( - repository, - job_id, - relay_id, - attempt.id(), - outcome, - evidence, - ) - .await + persist_outcome(repository, relay_id, outcome, &job, &attempt, evidence).await } async fn read_exact_event( @@ -239,19 +228,21 @@ async fn read_exact_event( async fn persist_outcome( repository: &MycStateRepository<'_>, - job_id: MycDeliveryJobId, relay_id: &MycDeliveryRelayId, - attempt_id: crate::MycDeliveryAttemptId, outcome: MycDeliveryAttemptOutcome, + job: &MycDeliveryJobRecord, + attempt: &MycDeliveryAttemptRecord, evidence: MycDeliveryExecutionEvidence, ) -> Result<MycDeliveryWorkerResult, MycDeliveryWorkerError> { + let retry_jitter = + delivery_retry_jitter(job, attempt.number(), outcome, evidence.retry_entropy)?; repository .record_delivery_attempt_outcome( - job_id, + job.id(), relay_id, - attempt_id, + attempt.id(), outcome, - evidence.retry_jitter, + retry_jitter, evidence.observed_at, ) .await @@ -259,6 +250,26 @@ async fn persist_outcome( .map_err(|_| worker_error(MycDeliveryWorkerErrorKind::State)) } +fn delivery_retry_jitter( + job: &MycDeliveryJobRecord, + attempt_number: u32, + outcome: MycDeliveryAttemptOutcome, + entropy: [u8; 8], +) -> Result<MycDeliveryRetryJitter, MycDeliveryWorkerError> { + if outcome == MycDeliveryAttemptOutcome::Delivered || attempt_number >= job.max_attempts() { + return MycDeliveryRetryJitter::new(0) + .map_err(|_| worker_error(MycDeliveryWorkerErrorKind::State)); + } + let exponent = attempt_number.saturating_sub(1).min(31); + let factor = 1_u64.checked_shl(exponent).unwrap_or(u64::MAX); + let maximum_delay = job + .initial_backoff_ms() + .saturating_mul(factor) + .min(job.maximum_backoff_ms()); + let jitter = u64::from_be_bytes(entropy) % (maximum_delay + 1); + MycDeliveryRetryJitter::new(jitter).map_err(|_| worker_error(MycDeliveryWorkerErrorKind::State)) +} + fn request_id(job_id: MycDeliveryJobId, attempt_id: crate::MycDeliveryAttemptId) -> String { let mut request = hex::encode(job_id.as_bytes()); request.push(':'); @@ -390,11 +401,55 @@ mod tests { submitted_at: MycDeliveryTimeUnixMs::new(RECEIVED_AT_MS + 4_001) .expect("submitted time"), observed_at: MycDeliveryTimeUnixMs::new(RECEIVED_AT_MS + 4_002).expect("observed time"), - retry_jitter: MycDeliveryRetryJitter::new(0).expect("jitter"), + retry_entropy: [0; 8], } } #[tokio::test] + async fn injected_entropy_uses_the_exact_full_jitter_cap() { + let root = tempfile::tempdir().expect("temporary root"); + let (host, job_id, _) = committed_response_host(root.path()).await; + let job = host + .repository() + .read_delivery_job(job_id) + .await + .expect("job read") + .expect("job"); + let first = delivery_retry_jitter( + &job, + 1, + MycDeliveryAttemptOutcome::TransportFailed, + u64::MAX.to_be_bytes(), + ) + .expect("first-attempt jitter"); + assert_eq!(first.get(), u64::MAX % (job.initial_backoff_ms() + 1)); + assert!(first.get() <= job.initial_backoff_ms()); + assert_eq!( + delivery_retry_jitter( + &job, + job.max_attempts(), + MycDeliveryAttemptOutcome::TransportFailed, + u64::MAX.to_be_bytes(), + ) + .expect("exhausted jitter") + .get(), + 0 + ); + assert_eq!( + delivery_retry_jitter( + &job, + 1, + MycDeliveryAttemptOutcome::Delivered, + u64::MAX.to_be_bytes(), + ) + .expect("terminal jitter") + .get(), + 0 + ); + host.close().await.expect("close"); + } + + #[tokio::test] async fn exact_bytes_are_submitted_only_after_durable_submitted_state() { let root = tempfile::tempdir().expect("temporary root"); let (host, job_id, exact_bytes) = committed_response_host(root.path()).await; diff --git a/src/lib.rs b/src/lib.rs @@ -36,6 +36,11 @@ mod provider_local_signer; mod provider_verification; mod runtime_context; mod runtime_foundation; +#[cfg(any(target_os = "linux", target_os = "macos"))] +mod runtime_graph; +#[cfg(any(target_os = "linux", target_os = "macos"))] +mod runtime_nip46; +mod runtime_signal; mod runtime_supervision; #[cfg(any(target_os = "linux", target_os = "macos"))] mod state_admin; @@ -118,7 +123,7 @@ pub use operations_v1::{ MycBoundOperationsServer, MycOperationsCancellationToken, MycOperationsError, MycOperationsErrorKind, MycOperationsServer, }; -pub use process_v1::execute_myc_cli_v1; +pub use process_v1::{execute_myc_cli_v1, execute_myc_cli_v1_with_signal_source}; pub use provider_contract::{ MYC_PROVIDER_CONCURRENCY_MAX, MYC_PROVIDER_CONTRACT_VERSION, MYC_PROVIDER_INPUT_MAX_BYTES, MYC_PROVIDER_OUTPUT_MAX_BYTES, MYC_PROVIDER_REQUEST_DEADLINE_MAX_MS, @@ -165,6 +170,7 @@ pub use runtime_foundation::{ MycRuntimeFoundationErrorKind, MycRuntimePrerequisite, MycRuntimeReadiness, MycRuntimeReadinessReason, open_myc_runtime_foundation, }; +pub use runtime_signal::{MycProcessSignal, MycProcessSignalFuture, MycProcessSignalSource}; pub use runtime_supervision::{ MYC_CRITICAL_TASK_MAX_COUNT, MYC_RUNTIME_SUPERVISION_CONTRACT_VERSION, MycCriticalTask, MycCriticalTaskError, MycRuntimeSupervisionError, MycRuntimeSupervisionErrorKind, diff --git a/src/main.rs b/src/main.rs @@ -4,10 +4,61 @@ use std::process::ExitCode; use myc::{MycLogRecord, MycProcessResult}; +struct MycOsSignalSource { + #[cfg(unix)] + interrupt: tokio::signal::unix::Signal, + #[cfg(unix)] + terminate: tokio::signal::unix::Signal, +} + +impl MycOsSignalSource { + fn new() -> Option<Self> { + #[cfg(unix)] + { + let interrupt = + tokio::signal::unix::signal(tokio::signal::unix::SignalKind::interrupt()).ok()?; + let terminate = + tokio::signal::unix::signal(tokio::signal::unix::SignalKind::terminate()).ok()?; + Some(Self { + interrupt, + terminate, + }) + } + #[cfg(not(unix))] + { + Some(Self {}) + } + } +} + +impl myc::MycProcessSignalSource for MycOsSignalSource { + fn next_signal(&mut self) -> myc::MycProcessSignalFuture<'_> { + #[cfg(unix)] + { + Box::pin(async move { + tokio::select! { + observed = self.interrupt.recv() => observed.map(|()| myc::MycProcessSignal::Interrupt), + observed = self.terminate.recv() => observed.map(|()| myc::MycProcessSignal::Terminate), + } + }) + } + #[cfg(not(unix))] + { + Box::pin(async move { + tokio::signal::ctrl_c() + .await + .ok() + .map(|()| myc::MycProcessSignal::Interrupt) + }) + } + } +} + fn main() -> ExitCode { match myc::parse_myc_cli_v1_from(std::env::args_os()) { Ok(invocation) => { - let result = myc::execute_myc_cli_v1(invocation); + let result = + myc::execute_myc_cli_v1_with_signal_source(invocation, MycOsSignalSource::new); eprintln!("{}", MycLogRecord::process_result(result)); result.exit_code() } diff --git a/src/nip46_replay.rs b/src/nip46_replay.rs @@ -110,6 +110,11 @@ impl MycReplayBoundNip46Request { self.event.client_public_key() } + #[must_use] + pub(crate) const fn encryption_context(&self) -> crate::MycNip46EncryptionContext { + self.event.encryption_context() + } + /// Returns the validated request identifier. #[must_use] pub const fn request_id(&self) -> &MycNip46RequestId { diff --git a/src/nip46_wave_080_a.rs b/src/nip46_wave_080_a.rs @@ -265,7 +265,7 @@ pub(crate) fn untrusted_response( } #[allow(clippy::too_many_arguments)] -async fn admit_connect( +pub(crate) async fn admit_connect( repository: &MycStateRepository<'_>, config: &MycConfigDocumentV1, client_seed: u8, diff --git a/src/nip46_work.rs b/src/nip46_work.rs @@ -217,6 +217,7 @@ impl fmt::Debug for MycDecryptedNip46Request { pub struct MycPreparedNip46Request { signer_request: MycSignerRequest, request: Request, + encryption_context: MycNip46EncryptionContext, } impl MycPreparedNip46Request { @@ -249,6 +250,7 @@ pub fn prepare_myc_nip46_request( nonce: MycSignerOperationNonce, received_at: MycRequestReceivedAtUnixMs, ) -> Result<MycPreparedNip46Request, MycNip46WorkError> { + let encryption_context = decrypted.replay.encryption_context(); let method = MycSignerRequestMethod::parse(decrypted.replay.method()) .ok_or_else(|| work_error(MycNip46WorkErrorKind::UnsupportedMethod))?; let message: RequestMessage = serde_json::from_slice(decrypted.replay.canonical_request()) @@ -273,6 +275,7 @@ pub fn prepare_myc_nip46_request( Ok(MycPreparedNip46Request { signer_request, request: message.request, + encryption_context, }) } @@ -298,6 +301,7 @@ pub struct MycNip46Work { request: MycSignerRequestRecord, connection: Option<MycConnectionRecord>, method: MycSignerRequestMethod, + encryption_context: MycNip46EncryptionContext, payload: Nip46WorkPayload, } @@ -320,6 +324,7 @@ impl MycNip46Work { request, connection: Some(connection), method, + encryption_context: MycNip46EncryptionContext::Nip44V2, payload: Nip46WorkPayload::Local, } } @@ -342,6 +347,10 @@ impl MycNip46Work { self.method } + pub(crate) const fn encryption_context(&self) -> MycNip46EncryptionContext { + self.encryption_context + } + /// Returns the work class without exposing protected parameters. #[must_use] pub const fn kind(&self) -> MycNip46WorkKind { @@ -424,6 +433,7 @@ pub fn prepare_myc_nip46_work( let MycPreparedNip46Request { signer_request, request, + encryption_context, } = prepared; if !record.matches_request(&signer_request) || observed_at.get() < record.received_at().get() { return Err(work_error(MycNip46WorkErrorKind::InvalidBinding)); @@ -448,6 +458,7 @@ pub fn prepare_myc_nip46_work( request: record, connection: None, method, + encryption_context, payload: Nip46WorkPayload::Connect { client: signer_request.client_public_key().clone(), permissions, @@ -562,6 +573,7 @@ pub fn prepare_myc_nip46_work( request: record, connection: Some(connection), method, + encryption_context, payload, }) } diff --git a/src/process_v1.rs b/src/process_v1.rs @@ -24,12 +24,13 @@ use crate::{ MYC_OPERATOR_CONTRACT_VERSION, MYC_PROVIDER_CONTRACT_VERSION, MYC_SIGNER_STATUS_CONTRACT_VERSION, MYC_STATE_SCHEMA_VERSION, MycBootstrapProfileV1, MycCliInvocationV1, MycCliOutputModeV1, MycCliPrimaryAuthorityV1, MycCommandV1, - MycConfigCommandV1, MycConfigDocumentV1, MycIdentityCommandArgsV1, MycIdentityCommandV1, - MycLocalSignerClient, MycProcessResult, MycProviderKind, MycProviderRole, MycRuntimeContext, - MycStateBackupArgsV1, MycStateCommandV1, MycStateMetadata, MycStateRestoreArgsV1, - RadrootsHostEnvironment, RadrootsPathResolver, RadrootsPlatform, finalize_myc_state_restore, - initialize_myc_config_document, initialize_myc_state, load_myc_config_candidate, - load_myc_config_document, open_myc_encrypted_identity, open_myc_state_inspection_from_config, + MycConfigCommandV1, MycConfigDocumentV1, MycConnectionCountsV1, MycIdentityCommandArgsV1, + MycIdentityCommandV1, MycLocalSignerClient, MycOutboxStatusV1, MycProcessResult, + MycProviderKind, MycProviderRole, MycRuntimeContext, MycStateBackupArgsV1, MycStateCommandV1, + MycStateMetadata, MycStateRestoreArgsV1, RadrootsHostEnvironment, RadrootsPathResolver, + RadrootsPlatform, finalize_myc_state_restore, initialize_myc_config_document, + initialize_myc_state, load_myc_config_candidate, load_myc_config_document, + open_myc_encrypted_identity, open_myc_state_inspection_from_config, open_myc_state_read_write_from_config, plan_myc_cli_v1, provision_myc_encrypted_identity, resolve_myc_runtime_context, resolve_myc_wrapping_credential, stage_myc_state_restore, verify_myc_state_backup, @@ -40,6 +41,12 @@ struct ProcessFailure(MycProcessResult); type ProcessResult<T> = Result<T, ProcessFailure>; +struct OfflineStateSnapshot { + value: Value, + connection_counts: MycConnectionCountsV1, + outbox: MycOutboxStatusV1, +} + /// Executes one admitted Myc invocation without reparsing process arguments. /// /// Result bytes are written to stdout only after the governed operation @@ -50,6 +57,56 @@ pub fn execute_myc_cli_v1(invocation: MycCliInvocationV1) -> MycProcessResult { execute(invocation).unwrap_or_else(|failure| failure.0) } +/// Executes one admitted invocation with a binary-owned process-signal source. +/// +/// The signal source factory is consulted only for `run` and only after the +/// governed Tokio runtime has entered. Non-daemon commands retain the exact +/// one-pass execution path used by [`execute_myc_cli_v1`]. +#[must_use] +pub fn execute_myc_cli_v1_with_signal_source<F, S>( + invocation: MycCliInvocationV1, + make_signal_source: F, +) -> MycProcessResult +where + F: FnOnce() -> Option<S>, + S: crate::MycProcessSignalSource + 'static, +{ + if !matches!(invocation.command(), MycCommandV1::Run) { + return execute_myc_cli_v1(invocation); + } + execute_run(invocation, make_signal_source).unwrap_or_else(|failure| failure.0) +} + +fn execute_run<F, S>( + invocation: MycCliInvocationV1, + make_signal_source: F, +) -> ProcessResult<MycProcessResult> +where + F: FnOnce() -> Option<S>, + S: crate::MycProcessSignalSource + 'static, +{ + let resolver = RadrootsPathResolver::new(RadrootsPlatform::current(), host_environment()); + let runtime = + resolve_myc_runtime_context(&resolver, &invocation).map_err(|_| input_failure())?; + let configuration = load_myc_config_document(&runtime).map_err(|_| input_failure())?; + let applied_at = migration_time()?; + let build = migration_build_identity()?; + let tokio = build_tokio_runtime(configuration.runtime_thread_limits())?; + tokio.block_on(async move { + let signals = make_signal_source().ok_or(ProcessFailure( + MycProcessResult::ServiceOrDependencyUnavailable, + ))?; + Ok(crate::runtime_graph::run_myc_daemon( + runtime, + configuration, + applied_at, + &build, + signals, + ) + .await) + }) +} + fn execute(invocation: MycCliInvocationV1) -> ProcessResult<MycProcessResult> { let plan = plan_myc_cli_v1(&invocation); if matches!( @@ -396,6 +453,15 @@ async fn offline_state_status( runtime: &MycRuntimeContext, configuration: &MycConfigDocumentV1, ) -> ProcessResult<Value> { + inspect_offline_state(runtime, configuration) + .await + .map(|snapshot| snapshot.value) +} + +async fn inspect_offline_state( + runtime: &MycRuntimeContext, + configuration: &MycConfigDocumentV1, +) -> ProcessResult<OfflineStateSnapshot> { let state = open_myc_state_inspection_from_config(runtime, configuration) .await .map_err(|_| state_failure())?; @@ -405,6 +471,16 @@ async fn offline_state_status( .current_configuration_generation() .await .map_err(|_| state_failure())?; + let connection_counts = state + .repository() + .read_runtime_connection_counts() + .await + .map_err(|_| state_failure())?; + let outbox = state + .repository() + .read_runtime_outbox_status() + .await + .map_err(|_| state_failure())?; let checked_at = integrity_time()?; let report = state .inspect_integrity(checked_at) @@ -417,14 +493,18 @@ async fn offline_state_status( .get(); let verified = report.sqlite() == IntegrityCheckOutcome::Verified && report.foreign_keys() == IntegrityCheckOutcome::Verified; - Ok(json!({ - "backup_eligible": verified, - "generation": generation, - "integrity": if verified { "verified" } else { "failed" }, - "reason_codes": if verified { json!([]) } else { json!(["database_integrity_failed"]) }, - "schema_version": schema, - "writer_lock": "free", - })) + Ok(OfflineStateSnapshot { + value: json!({ + "backup_eligible": verified, + "generation": generation, + "integrity": if verified { "verified" } else { "failed" }, + "reason_codes": if verified { json!([]) } else { json!(["database_integrity_failed"]) }, + "schema_version": schema, + "writer_lock": "free", + }), + connection_counts, + outbox, + }) } .await; let closed = state.close().await.map_err(|_| state_failure()); @@ -435,7 +515,8 @@ async fn offline_service_status( runtime: &MycRuntimeContext, configuration: &MycConfigDocumentV1, ) -> ProcessResult<Value> { - let state = offline_state_status(runtime, configuration).await?; + let snapshot = inspect_offline_state(runtime, configuration).await?; + let state = &snapshot.value; let generation = state .get("generation") .and_then(Value::as_u64) @@ -492,9 +573,9 @@ async fn offline_service_status( "contract_version": MYC_SIGNER_STATUS_CONTRACT_VERSION, "instance": runtime.context().instance().as_str(), "myc": { - "connection_counts": {}, + "connection_counts": snapshot.connection_counts, "discovery": unavailable(discovery_configured), - "outbox": {"pending": 0, "unknown": 0}, + "outbox": snapshot.outbox, "transport": unavailable(true), "user": unavailable(true), }, diff --git a/src/process_v1_unsupported.rs b/src/process_v1_unsupported.rs @@ -12,3 +12,16 @@ pub fn execute_myc_cli_v1(invocation: MycCliInvocationV1) -> MycProcessResult { MycProcessResult::InputOrConfiguration } } + +/// Rejects execution without consulting a signal source on unsupported hosts. +#[must_use] +pub fn execute_myc_cli_v1_with_signal_source<F, S>( + invocation: MycCliInvocationV1, + _make_signal_source: F, +) -> MycProcessResult +where + F: FnOnce() -> Option<S>, + S: crate::MycProcessSignalSource + 'static, +{ + execute_myc_cli_v1(invocation) +} diff --git a/src/provider_contract.rs b/src/provider_contract.rs @@ -810,6 +810,49 @@ impl MycProviderOperationInput { } } + fn owned(&self) -> Self { + let kind = match &self.kind { + ProviderOperationInputKind::Describe => ProviderOperationInputKind::Describe, + ProviderOperationInputKind::PublicIdentity => { + ProviderOperationInputKind::PublicIdentity + } + ProviderOperationInputKind::SignEvent(bytes) => { + ProviderOperationInputKind::SignEvent(bytes.clone()) + } + ProviderOperationInputKind::Nip04Encrypt { peer, plaintext } => { + ProviderOperationInputKind::Nip04Encrypt { + peer: peer.clone(), + plaintext: plaintext.clone(), + } + } + ProviderOperationInputKind::Nip04Decrypt { peer, ciphertext } => { + ProviderOperationInputKind::Nip04Decrypt { + peer: peer.clone(), + ciphertext: ciphertext.clone(), + } + } + ProviderOperationInputKind::Nip44Encrypt { + peer, + version, + plaintext, + } => ProviderOperationInputKind::Nip44Encrypt { + peer: peer.clone(), + version: *version, + plaintext: plaintext.clone(), + }, + ProviderOperationInputKind::Nip44Decrypt { + peer, + version, + ciphertext, + } => ProviderOperationInputKind::Nip44Decrypt { + peer: peer.clone(), + version: *version, + ciphertext: ciphertext.clone(), + }, + }; + Self { kind } + } + #[must_use] pub const fn nip44_version(&self) -> Option<MycProviderNip44Version> { match &self.kind { @@ -925,6 +968,19 @@ impl MycProviderOperation { &self.input } + pub(crate) fn owned_for_runtime(&self) -> Self { + Self { + role: self.role, + instance: self.instance, + provider: self.provider, + operation_id: self.operation_id, + correlation_id: self.correlation_id, + deadline: self.deadline, + expected_identity: self.expected_identity.clone(), + input: self.input.owned(), + } + } + pub(crate) fn binding_digest(&self) -> [u8; 32] { let mut hasher = Sha256::new(); hasher.update(b"radroots.myc.provider.operation_binding.v1\0"); diff --git a/src/runtime_graph.rs b/src/runtime_graph.rs @@ -0,0 +1,2761 @@ +//! Final native Myc daemon graph, readiness, reconnect, and bounded shutdown. + +use core::{fmt, future::Future, future::pending, pin::Pin, time::Duration}; +use std::{ + collections::VecDeque, + error::Error, + path::PathBuf, + sync::{ + Arc, + atomic::{AtomicBool, Ordering}, + }, +}; + +use base64::{Engine as _, engine::general_purpose::URL_SAFE_NO_PAD}; +use radroots_service_host::{ + EntropySource, GracefulShutdown, HostError, HostErrorKind, MetricValue, MonotonicClock, + ProcessSignal, ProcessSignalAdapter, ProcessSignalFuture, ProcessSignalSource, + ShutdownDisposition, ShutdownPhase, ShutdownPhaseFuture, ShutdownPhaseHandler, + SupervisedTaskExitStatus, SystemEntropy, SystemMonotonicClock, SystemWallClock, + TaskClassification, TaskMetadata, TaskName, TaskSupervisor, WallClock, +}; +use radroots_service_sqlite::{ + BackupCreatedAtUnixMs, MigrationAppliedAtUnixSeconds, MigrationBuildIdentity, +}; +use radroots_transport::{BoxSubscription, SubscriptionNext, source::SubscriptionCheckpoint}; +use serde_json::{Value, json}; +use sha2::{Digest, Sha256}; +use tokio::sync::mpsc; + +use crate::delivery_worker::{MycDeliveryExecutionEvidence, MycDeliveryWorker}; +use crate::provider_executor::MycProviderExecutor; +use crate::runtime_nip46::{ + MycNip46DispatchErrorKind, MycRuntimeNip46AdmissionEvidence, MycRuntimeNip46Coordinator, +}; +use crate::state_admin::AdminJournalOperationError; +use crate::state_connection::{ + MycAdminConnectionAction, apply_admin_challenge_authorize, apply_admin_challenge_require, + apply_admin_connection_mutation, +}; +use crate::state_discovery::{apply_admin_discovery_publish, prepare_discovery_signing_bytes}; +use crate::state_governance::MycAdminAuditQuery; +use crate::transport_nostr_adapter::MycNostrIngressAdapter; +use crate::{ + MycAdminCancellationToken, MycAdminFuture, MycAdminHandler, MycAdminHandlerError, + MycAdminHandlerErrorKind, MycAdminOperationAdmission, MycAdminOperationErrorKind, + MycAdminOperationJournalPolicy, MycAdminOperationTimeUnixMs, MycAdminRequestDocument, + MycAdminResponseDocument, MycAdminRoute, MycAdminServer, MycAuditCorrelationId, MycAuditKind, + MycAuditOutcome, MycAuditPageLimit, MycAuthorizationChallengeId, + MycAuthorizationChallengeNonce, MycAuthorizationChallengeRequest, + MycAuthorizationChallengeState, MycAuthorizationChallengeUrl, MycBootstrapProfileV1, + MycBoundAdminServer, MycBoundOperationsServer, MycConfigDocumentV1, MycConnectionCountsV1, + MycConnectionId, MycConnectionPermission, MycConnectionPermissionSet, + MycConnectionPolicyGeneration, MycConnectionStatus, MycConnectionTimeUnixMs, + MycDeliveryAttemptNonce, MycDeliveryRecoveryEntropy, MycDeliveryTimeUnixMs, + MycDiscoveryCommitRequest, MycIdentityHealthV1, MycOperationsCancellationToken, + MycOperationsServer, MycOutboxStatusV1, MycPersistenceHealthV1, MycPersistenceStatusV1, + MycProcessResult, MycProcessSignal, MycProcessSignalSource, MycProviderCorrelationId, + MycProviderDeadlineUnixMs, MycProviderOperation, MycProviderOperationId, + MycProviderOperationInput, MycProviderResponseObservedAtUnixMs, MycProviderRole, + MycProviderStatusV1, MycRateLimitClass, MycRelayTransportStatusV1, MycRuntimeContext, + MycServicePhase, MycSignerOperationId, MycStateHost, MycStatusBuildInfoV1, MycStatusBuildMode, + MycStatusCommonV1, MycStatusConfigurationIdentityV1, MycStatusConfigurationSource, + MycStatusObservationV1, MycStatusPublisher, MycStatusReader, MycStatusReasonCode, + MycStatusReasonCodes, MycTaskCancellation, MycTransportHealthV1, myc_status_cache, + open_myc_state_read_write_from_config, +}; + +const TASK_ADMIN_SERVER: &str = "admin_server"; +const TASK_OPERATIONS_SERVER: &str = "operations_server"; +const TASK_RELAY_INGRESS: &str = "relay_ingress"; +const TASK_PROVIDER_DISPATCH: &str = "provider_dispatch"; +const TASK_DELIVERY_OUTBOX: &str = "delivery_outbox"; + +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +enum MycDaemonErrorKind { + State, + Provider, + Relay, + Admin, + Operations, + Runtime, +} + +struct MycDaemonError { + kind: MycDaemonErrorKind, +} + +impl MycDaemonError { + const fn new(kind: MycDaemonErrorKind) -> Self { + Self { kind } + } + + const fn process_result(&self) -> MycProcessResult { + match self.kind { + MycDaemonErrorKind::State => MycProcessResult::StateOrIdentityUnavailable, + MycDaemonErrorKind::Provider + | MycDaemonErrorKind::Relay + | MycDaemonErrorKind::Admin + | MycDaemonErrorKind::Operations => MycProcessResult::ServiceOrDependencyUnavailable, + MycDaemonErrorKind::Runtime => MycProcessResult::UnexpectedInternal, + } + } +} + +impl fmt::Debug for MycDaemonError { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter + .debug_struct("MycDaemonError") + .field("kind", &self.kind) + .finish() + } +} + +impl fmt::Display for MycDaemonError { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter.write_str("Myc daemon failed") + } +} + +impl Error for MycDaemonError {} + +struct HostSignalSource<S> { + inner: S, +} + +impl<S> HostSignalSource<S> { + const fn new(inner: S) -> Self { + Self { inner } + } +} + +impl<S> ProcessSignalSource for HostSignalSource<S> +where + S: MycProcessSignalSource, +{ + fn next_signal(&mut self) -> ProcessSignalFuture<'_> { + Box::pin(async move { + self.inner.next_signal().await.map(|signal| match signal { + MycProcessSignal::Interrupt => ProcessSignal::Interrupt, + #[cfg(unix)] + MycProcessSignal::Terminate => ProcessSignal::Terminate, + }) + }) + } +} + +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +enum HealthEvent { + Providers(bool), + RequiredRelays(bool), + StateChanged, +} + +struct RuntimeHealth { + providers: bool, + required_relays: bool, + operations: bool, +} + +struct RuntimeIngressItem { + raw_event: Box<[u8]>, + relay_id: crate::MycRateRelayId, + observed_at_unix_ms: u64, +} + +struct RuntimeIngressSlot { + index: usize, + required: bool, + checkpoints: Vec<SubscriptionCheckpoint>, + state: RuntimeIngressSlotState, + retry_delay_ms: u64, +} + +enum RuntimeIngressSlotState { + Active(BoxSubscription), + Waiting, + Connecting( + Pin< + Box< + dyn Future< + Output = Result< + BoxSubscription, + crate::transport_nostr_adapter::MycRelayAdapterError, + >, + > + Send, + >, + >, + ), +} + +enum RuntimeIngressAction { + Next(Result<SubscriptionNext, radroots_transport::Error>), + Retry, + Connected(Result<BoxSubscription, crate::transport_nostr_adapter::MycRelayAdapterError>), +} + +impl RuntimeHealth { + const fn ready(operations: bool) -> Self { + Self { + providers: true, + required_relays: true, + operations, + } + } + + fn observe(&mut self, event: HealthEvent) -> bool { + match event { + HealthEvent::Providers(value) => { + let changed = self.providers != value; + self.providers = value; + changed + } + HealthEvent::RequiredRelays(value) => { + let changed = self.required_relays != value; + self.required_relays = value; + changed + } + HealthEvent::StateChanged => true, + } + } +} + +struct RuntimeStatusContext { + runtime: MycRuntimeContext, + configuration: Arc<MycConfigDocumentV1>, + state: Arc<MycStateHost>, + clock: SystemMonotonicClock, + generation: u64, + connection_counts: MycConnectionCountsV1, + outbox: MycOutboxStatusV1, +} + +impl RuntimeStatusContext { + async fn refresh_state(&mut self) -> Result<(), MycDaemonError> { + let connection_counts = self + .state + .repository() + .read_runtime_connection_counts() + .await + .map_err(|_| MycDaemonError::new(MycDaemonErrorKind::State))?; + let outbox = self + .state + .repository() + .read_runtime_outbox_status() + .await + .map_err(|_| MycDaemonError::new(MycDaemonErrorKind::State))?; + self.connection_counts = connection_counts; + self.outbox = outbox; + Ok(()) + } + + fn observation( + &self, + phase: MycServicePhase, + health: &RuntimeHealth, + ) -> Result<MycStatusObservationV1, MycDaemonError> { + let ready = matches!(phase, MycServicePhase::Ready | MycServicePhase::Degraded) + && health.providers + && health.required_relays; + let mut reasons = Vec::new(); + if !health.providers { + reasons.push(MycStatusReasonCode::SignerProviderUnavailable); + } + if !health.required_relays { + reasons.push(MycStatusReasonCode::RequiredRelayUnavailable); + reasons.push(MycStatusReasonCode::SubscriberNotActive); + } + if !health.operations { + reasons.push(MycStatusReasonCode::OperationsListenerFailed); + } + if phase == MycServicePhase::Stopping { + reasons.push(MycStatusReasonCode::ShutdownInProgress); + } + let reasons = MycStatusReasonCodes::new(reasons) + .map_err(|_| MycDaemonError::new(MycDaemonErrorKind::Runtime))?; + let build = runtime_build_info()?; + let configuration = MycStatusConfigurationIdentityV1::new( + hex::encode(self.state.metadata().configuration_digest().as_bytes()), + if self.runtime.profile() == MycBootstrapProfileV1::RepoLocal { + MycStatusConfigurationSource::DerivedRepoLocal + } else { + MycStatusConfigurationSource::ExplicitConfig + }, + ) + .map_err(|_| MycDaemonError::new(MycDaemonErrorKind::Runtime))?; + let persistence = MycPersistenceStatusV1::new( + MycPersistenceHealthV1::Ready, + self.state + .metadata() + .initial_database_metadata() + .state_schema_version() + .get(), + self.generation, + crate::MycIntegrityStateV1::Verified, + MycStatusReasonCodes::empty(), + ) + .map_err(|_| MycDaemonError::new(MycDaemonErrorKind::Runtime))?; + let common = MycStatusCommonV1::new( + phase, + ready, + reasons.clone(), + u64::try_from( + self.clock + .now_monotonic() + .duration_since_origin() + .as_millis(), + ) + .unwrap_or(u64::MAX), + build, + configuration, + persistence, + ) + .map_err(|_| MycDaemonError::new(MycDaemonErrorKind::Runtime))?; + let identity = |configured: bool| { + MycIdentityHealthV1::new(configured, configured && health.providers, reasons.clone()) + }; + let provider = MycProviderStatusV1::new( + identity(true).map_err(|_| MycDaemonError::new(MycDaemonErrorKind::Runtime))?, + identity(true).map_err(|_| MycDaemonError::new(MycDaemonErrorKind::Runtime))?, + identity( + self.configuration + .provider_contract() + .binding(MycProviderRole::Discovery) + .is_some(), + ) + .map_err(|_| MycDaemonError::new(MycDaemonErrorKind::Runtime))?, + reasons.clone(), + ) + .map_err(|_| MycDaemonError::new(MycDaemonErrorKind::Runtime))?; + let transport = MycRelayTransportStatusV1::new( + if health.required_relays { + MycTransportHealthV1::Ready + } else { + MycTransportHealthV1::Unavailable + }, + health.required_relays, + if health.required_relays { + u64::try_from(self.configuration.relay_count()).unwrap_or(u64::MAX) + } else { + 0 + }, + reasons, + ) + .map_err(|_| MycDaemonError::new(MycDaemonErrorKind::Runtime))?; + Ok(MycStatusObservationV1::new( + common, + provider, + transport, + self.connection_counts, + self.outbox, + )) + } +} + +struct RuntimeAdminHandler { + state: Arc<MycStateHost>, + configuration: Arc<MycConfigDocumentV1>, + providers: Arc<MycProviderExecutor>, + status: MycStatusReader, + accepting_mutations: Arc<AtomicBool>, + cursor_key: [u8; 32], + health: mpsc::Sender<HealthEvent>, +} + +impl RuntimeAdminHandler { + async fn handle_inner( + &self, + request: MycAdminRequestDocument, + ) -> Result<MycAdminResponseDocument, MycAdminHandlerError> { + if request.route().is_mutation() && !self.accepting_mutations.load(Ordering::Acquire) { + return Err(handler_error(MycAdminHandlerErrorKind::Unavailable)); + } + match request.route() { + MycAdminRoute::Status => MycAdminResponseDocument::from_canonical_bytes( + request.route(), + self.status.snapshot().detailed_status_json(), + ) + .map_err(|_| handler_error(MycAdminHandlerErrorKind::Internal)), + MycAdminRoute::EffectiveConfig => MycAdminResponseDocument::from_canonical_bytes( + request.route(), + self.configuration.effective().canonical_json().as_bytes(), + ) + .map_err(|_| handler_error(MycAdminHandlerErrorKind::Internal)), + MycAdminRoute::IdentityStatus | MycAdminRoute::IdentityPublic => { + self.identity_response(request).await + } + MycAdminRoute::StateStatus => self.state_status(request).await, + MycAdminRoute::StateBackup => self.backup(request).await, + MycAdminRoute::MetricsSnapshot => self.metrics_snapshot(request), + MycAdminRoute::ConnectionsList => self.connections(request).await, + MycAdminRoute::AuditEvents => self.audit_events(request).await, + MycAdminRoute::AuditSummary => self.audit_summary(request).await, + MycAdminRoute::DiscoveryDesired => self.discovery_desired(request).await, + MycAdminRoute::ConnectionApprove + | MycAdminRoute::ConnectionReject + | MycAdminRoute::ConnectionRevoke => self.connection_mutation(request).await, + MycAdminRoute::ChallengeRequire => self.challenge_require(request).await, + MycAdminRoute::ChallengeAuthorize => self.challenge_authorize(request).await, + MycAdminRoute::DiscoveryRender + | MycAdminRoute::DiscoveryRefresh + | MycAdminRoute::DiscoveryPublish => self.discovery_mutation(request).await, + } + } + + async fn identity_response( + &self, + request: MycAdminRequestDocument, + ) -> Result<MycAdminResponseDocument, MycAdminHandlerError> { + let model: Value = serde_json::from_slice(request.model_bytes()) + .map_err(|_| handler_error(MycAdminHandlerErrorKind::Internal))?; + let role_name = model + .pointer("/role") + .and_then(Value::as_str) + .ok_or_else(|| handler_error(MycAdminHandlerErrorKind::Internal))?; + let role = match role_name { + "transport" => MycProviderRole::Transport, + "user" => MycProviderRole::User, + "discovery" => MycProviderRole::Discovery, + _ => return Err(handler_error(MycAdminHandlerErrorKind::Internal)), + }; + let binding = self.configuration.provider_contract().binding(role); + let expected = match role { + MycProviderRole::Transport => { + Some(self.state.metadata().expected_identities().transport()) + } + MycProviderRole::User => Some(self.state.metadata().expected_identities().user()), + MycProviderRole::Discovery => self.state.metadata().expected_identities().discovery(), + }; + let generation = u64::from( + self.state + .repository() + .current_configuration_generation() + .await + .map_err(|_| handler_error(MycAdminHandlerErrorKind::Internal))?, + ); + let value = if request.route() == MycAdminRoute::IdentityPublic { + let expected = + expected.ok_or_else(|| handler_error(MycAdminHandlerErrorKind::NotFound))?; + json!({"generation":generation,"public_key":expected.as_hex(),"role":role_name}) + } else { + let mut value = json!({ + "available": binding.is_some(), + "configured": binding.is_some(), + "generation": generation, + "reason_codes": [], + "role": role_name, + }); + if let Some(binding) = binding { + value["provider"] = Value::String(binding.kind().as_str().to_owned()); + } + if let Some(expected) = expected { + value["public_key"] = Value::String(expected.as_hex().to_owned()); + } + value + }; + response(request.route(), value) + } + + async fn state_status( + &self, + request: MycAdminRequestDocument, + ) -> Result<MycAdminResponseDocument, MycAdminHandlerError> { + let generation = self + .state + .repository() + .current_configuration_generation() + .await + .map_err(|_| handler_error(MycAdminHandlerErrorKind::Internal))?; + response( + request.route(), + json!({ + "backup_eligible": true, + "generation": u64::from(generation), + "integrity": "verified", + "reason_codes": [], + "schema_version": self.state.metadata().initial_database_metadata().state_schema_version().get(), + "writer_lock": "held_by_daemon", + }), + ) + } + + async fn backup( + &self, + request: MycAdminRequestDocument, + ) -> Result<MycAdminResponseDocument, MycAdminHandlerError> { + let now_seconds = SystemWallClock + .now_utc() + .map_err(|_| handler_error(MycAdminHandlerErrorKind::Internal))? + .get(); + let now_millis = now_seconds + .checked_mul(1_000) + .ok_or_else(|| handler_error(MycAdminHandlerErrorKind::Internal))?; + let prepared_at = MycAdminOperationTimeUnixMs::new(now_millis) + .map_err(|_| handler_error(MycAdminHandlerErrorKind::Internal))?; + let model: Value = serde_json::from_slice(request.model_bytes()) + .map_err(|_| handler_error(MycAdminHandlerErrorKind::Internal))?; + let expected_generation = model + .pointer("/expected_generation") + .and_then(Value::as_u64) + .ok_or_else(|| handler_error(MycAdminHandlerErrorKind::Internal))?; + let actual_generation = u64::from( + self.state + .repository() + .current_configuration_generation() + .await + .map_err(|_| handler_error(MycAdminHandlerErrorKind::Internal))?, + ); + if expected_generation != actual_generation { + return Err(handler_error(MycAdminHandlerErrorKind::Conflict)); + } + let prepared = match self + .state + .repository() + .prepare_admin_operation(&request, prepared_at) + .await + .map_err(map_admin_journal_error)? + { + MycAdminOperationAdmission::ExactReplay(response) => return Ok(response), + MycAdminOperationAdmission::Prepared(prepared) => prepared, + }; + let target = model + .pointer("/target_path") + .and_then(Value::as_str) + .map(PathBuf::from) + .ok_or_else(|| handler_error(MycAdminHandlerErrorKind::Internal))?; + let manifest = self + .state + .capture_online_backup( + &target, + BackupCreatedAtUnixMs::new(now_millis) + .map_err(|_| handler_error(MycAdminHandlerErrorKind::Internal))?, + ) + .await + .map_err(|_| handler_error(MycAdminHandlerErrorKind::Unavailable))?; + let completed_at_seconds = SystemWallClock + .now_utc() + .map_err(|_| handler_error(MycAdminHandlerErrorKind::Internal))? + .get(); + let completed_at = MycAdminOperationTimeUnixMs::new( + completed_at_seconds + .checked_mul(1_000) + .ok_or_else(|| handler_error(MycAdminHandlerErrorKind::Internal))?, + ) + .map_err(|_| handler_error(MycAdminHandlerErrorKind::Internal))?; + let operation_id = request + .operation_id() + .ok_or_else(|| handler_error(MycAdminHandlerErrorKind::Internal))?; + let response = response( + request.route(), + json!({ + "completed_at_utc": completed_at_seconds, + "manifest_digest": hex::encode(manifest.digest().as_bytes()), + "operation_id": operation_id, + "snapshot_generation": actual_generation, + }), + )?; + self.state + .repository() + .complete_admin_operation( + &prepared, + &response, + completed_at, + MycAdminOperationJournalPolicy::seven_days(), + ) + .await + .map_err(map_admin_journal_error)?; + Ok(response) + } + + fn metrics_snapshot( + &self, + request: MycAdminRequestDocument, + ) -> Result<MycAdminResponseDocument, MycAdminHandlerError> { + let captured_at = SystemWallClock + .now_utc() + .map_err(|_| handler_error(MycAdminHandlerErrorKind::Internal))? + .get(); + let snapshot = self.status.operations_cache().snapshot(); + let mut metrics = serde_json::Map::new(); + for sample in snapshot.metrics().samples() { + let mut key = sample.name().as_str().to_owned(); + for label in sample.labels() { + key.push('_'); + key.push_str(label.value()); + } + if !safe_metric_key(&key) { + return Err(handler_error(MycAdminHandlerErrorKind::Internal)); + } + let value = match sample.value() { + MetricValue::Counter(value) => value, + MetricValue::Gauge(value) => u64::try_from(value) + .map_err(|_| handler_error(MycAdminHandlerErrorKind::Internal))?, + }; + if metrics.insert(key, Value::from(value)).is_some() { + return Err(handler_error(MycAdminHandlerErrorKind::Internal)); + } + } + response( + request.route(), + json!({"captured_at_utc":captured_at,"metrics":metrics}), + ) + } + + async fn connections( + &self, + request: MycAdminRequestDocument, + ) -> Result<MycAdminResponseDocument, MycAdminHandlerError> { + let model = request_model(&request)?; + let limit = model_u64(&model, "/limit") + .and_then(|value| u16::try_from(value).ok()) + .ok_or_else(|| handler_error(MycAdminHandlerErrorKind::Internal))?; + let status = model + .pointer("/state") + .and_then(Value::as_str) + .map(|value| match value { + "pending" => Ok(MycConnectionStatus::Pending), + "approved" => Ok(MycConnectionStatus::Active), + "rejected" => Ok(MycConnectionStatus::Denied), + "revoked" => Ok(MycConnectionStatus::Expired), + _ => Err(handler_error(MycAdminHandlerErrorKind::Internal)), + }) + .transpose()?; + let cursor = model + .pointer("/cursor") + .and_then(Value::as_str) + .map(|value| decode_connection_cursor(value, status, &self.cursor_key)) + .transpose()?; + let snapshot = match cursor.as_ref() { + Some(cursor) => cursor.snapshot, + None => connection_time_now_for_admin()?, + }; + let before = cursor.map(|cursor| (cursor.before, cursor.id)); + let page = self + .state + .repository() + .read_admin_connection_page(limit, status, snapshot, before) + .await + .map_err(|_| handler_error(MycAdminHandlerErrorKind::Internal))?; + let items = page + .items() + .iter() + .map(|record| { + json!({ + "client_public_key": record.client_public_key().as_hex(), + "connection_id": hex::encode(record.id().as_bytes()), + "created_at_utc": record.created_at().get() / 1_000, + "generation": record.policy_generation().get(), + "permissions": record.admin_permissions(), + "state": record.status().admin_state(), + "updated_at_utc": record.updated_at().get() / 1_000, + }) + }) + .collect::<Vec<_>>(); + let mut value = json!({ + "items": items, + "snapshot_generation": snapshot.get(), + }); + if let Some((before, id)) = page.next() { + value["next_cursor"] = Value::String(encode_connection_cursor( + snapshot, + before, + id, + status, + &self.cursor_key, + )); + } + response(request.route(), value) + } + + async fn audit_events( + &self, + request: MycAdminRequestDocument, + ) -> Result<MycAdminResponseDocument, MycAdminHandlerError> { + let model = request_model(&request)?; + let limit = model_u64(&model, "/limit") + .and_then(|value| u16::try_from(value).ok()) + .and_then(|value| MycAuditPageLimit::new(value).ok()) + .ok_or_else(|| handler_error(MycAdminHandlerErrorKind::Internal))?; + let from_ms = optional_seconds_as_millis(&model, "/from_utc", false)?; + let to_ms = optional_seconds_as_millis(&model, "/to_utc", true)?; + let kind = model + .pointer("/kind") + .and_then(Value::as_str) + .map(|value| { + MycAuditKind::parse(value) + .ok_or_else(|| handler_error(MycAdminHandlerErrorKind::Internal)) + }) + .transpose()?; + let outcome = model + .pointer("/outcome") + .and_then(Value::as_str) + .map(|value| { + MycAuditOutcome::parse(value) + .ok_or_else(|| handler_error(MycAdminHandlerErrorKind::Internal)) + }) + .transpose()?; + let query_digest = audit_query_digest(from_ms, to_ms, kind, outcome); + let cursor = model + .pointer("/cursor") + .and_then(Value::as_str) + .map(|value| decode_audit_cursor(value, query_digest, &self.cursor_key)) + .transpose()?; + let page = self + .state + .repository() + .read_admin_audit_page( + MycAdminAuditQuery::new( + limit, + cursor.as_ref().map(|value| value.snapshot), + cursor.as_ref().map(|value| value.before), + from_ms, + to_ms, + kind, + outcome, + ) + .map_err(|_| handler_error(MycAdminHandlerErrorKind::Internal))?, + ) + .await + .map_err(|_| handler_error(MycAdminHandlerErrorKind::Internal))?; + let items = page + .items() + .iter() + .map(|record| { + json!({ + "audit_id": format!("audit-{}", record.sequence()), + "correlation_id": hex::encode(record.correlation_id().as_bytes()), + "kind": record.kind().as_str(), + "occurred_at_utc": record.occurred_at().get() / 1_000, + "outcome": record.outcome().as_str(), + "reason_code": record.reason().as_str(), + }) + }) + .collect::<Vec<_>>(); + let mut value = json!({ + "items": items, + "snapshot_generation": page.snapshot_sequence(), + }); + if let Some(before) = page.next_before_sequence() { + value["next_cursor"] = Value::String(encode_audit_cursor( + page.snapshot_sequence(), + before, + query_digest, + &self.cursor_key, + )); + } + response(request.route(), value) + } + + async fn audit_summary( + &self, + request: MycAdminRequestDocument, + ) -> Result<MycAdminResponseDocument, MycAdminHandlerError> { + let model = request_model(&request)?; + let from = model_u64(&model, "/from_utc") + .ok_or_else(|| handler_error(MycAdminHandlerErrorKind::Internal))?; + let to = model_u64(&model, "/to_utc") + .ok_or_else(|| handler_error(MycAdminHandlerErrorKind::Internal))?; + let from_ms = seconds_as_millis(from, false)?; + let to_ms = seconds_as_millis(to, true)?; + let counts = self + .state + .repository() + .read_admin_audit_summary(from_ms, to_ms) + .await + .map_err(|_| handler_error(MycAdminHandlerErrorKind::Internal))? + .iter() + .map(|(key, value)| (key.clone(), Value::from(*value))) + .collect::<serde_json::Map<_, _>>(); + response( + request.route(), + json!({ + "counts": counts, + "from_utc": from, + "to_utc": to, + }), + ) + } + + async fn connection_mutation( + &self, + request: MycAdminRequestDocument, + ) -> Result<MycAdminResponseDocument, MycAdminHandlerError> { + let route = request.route(); + let model = request_model(&request)?; + let operation_id = request_operation_id(&request)?.to_owned(); + let connection_id = + digest_id_parameter(&request, "connection_id").map(MycConnectionId::from_bytes)?; + let generation = model_u64(&model, "/expected_generation") + .and_then(|value| MycConnectionPolicyGeneration::new(value).ok()) + .ok_or_else(|| handler_error(MycAdminHandlerErrorKind::Internal))?; + let observed_at = connection_time_now_for_admin()?; + let correlation = admin_audit_correlation(&request); + let action = match route { + MycAdminRoute::ConnectionApprove => { + let permissions = model + .pointer("/permissions") + .and_then(Value::as_str) + .ok_or_else(|| handler_error(MycAdminHandlerErrorKind::Internal))?; + let permissions = if permissions.is_empty() { + Vec::new() + } else { + permissions + .split(',') + .map(|value| { + MycConnectionPermission::parse(value) + .ok_or_else(|| handler_error(MycAdminHandlerErrorKind::Internal)) + }) + .collect::<Result<Vec<_>, _>>()? + }; + let permissions = MycConnectionPermissionSet::new(&permissions) + .map_err(|_| handler_error(MycAdminHandlerErrorKind::Internal))?; + let authorized_until = self + .configuration + .normalized() + .pointer("/policy/challenges/enabled") + .and_then(Value::as_bool) + .unwrap_or(false) + .then(|| { + let lifetime = configuration_integer( + &self.configuration, + "/policy/challenges/authorized_lifetime_ms", + )?; + let until = observed_at + .get() + .checked_add(lifetime) + .ok_or_else(|| MycDaemonError::new(MycDaemonErrorKind::Runtime))?; + MycConnectionTimeUnixMs::new(until) + .map_err(|_| MycDaemonError::new(MycDaemonErrorKind::Runtime)) + }) + .transpose() + .map_err(|_| handler_error(MycAdminHandlerErrorKind::Internal))?; + MycAdminConnectionAction::Approve { + permissions, + authorized_until, + } + } + MycAdminRoute::ConnectionReject => MycAdminConnectionAction::Reject, + MycAdminRoute::ConnectionRevoke => MycAdminConnectionAction::Revoke, + _ => return Err(handler_error(MycAdminHandlerErrorKind::Internal)), + }; + let completed_at = admin_operation_time(observed_at)?; + self.state + .repository() + .execute_database_admin_operation( + &request, + completed_at, + MycAdminOperationJournalPolicy::seven_days(), + move |transaction| { + Box::pin(async move { + let (before, after) = apply_admin_connection_mutation( + transaction, + connection_id, + generation, + observed_at, + correlation, + action, + ) + .await?; + response( + route, + json!({ + "connection_id": hex::encode(after.id().as_bytes()), + "current_state": after.status().admin_state(), + "generation": after.policy_generation().get(), + "operation_id": operation_id, + "previous_state": before.status().admin_state(), + }), + ) + .map_err(|_| AdminJournalOperationError::Binding) + }) + }, + ) + .await + .map_err(map_admin_journal_error) + } + + async fn challenge_require( + &self, + request: MycAdminRequestDocument, + ) -> Result<MycAdminResponseDocument, MycAdminHandlerError> { + let route = request.route(); + let model = request_model(&request)?; + let connection_id = model + .pointer("/connection_id") + .and_then(Value::as_str) + .ok_or_else(|| handler_error(MycAdminHandlerErrorKind::Internal)) + .and_then(digest_id) + .map(MycConnectionId::from_bytes)?; + let operation_id = model + .pointer("/request_identity") + .and_then(Value::as_str) + .ok_or_else(|| handler_error(MycAdminHandlerErrorKind::Internal)) + .and_then(digest_id) + .map(MycSignerOperationId::from_persisted)?; + let generation = model_u64(&model, "/expected_policy_generation") + .and_then(|value| MycConnectionPolicyGeneration::new(value).ok()) + .ok_or_else(|| handler_error(MycAdminHandlerErrorKind::Internal))?; + let issued_at = connection_time_now_for_admin()?; + let lifetime = configuration_integer( + &self.configuration, + "/policy/challenges/pending_lifetime_ms", + ) + .map_err(|_| handler_error(MycAdminHandlerErrorKind::Internal))?; + let expires_at = MycConnectionTimeUnixMs::new( + issued_at + .get() + .checked_add(lifetime) + .ok_or_else(|| handler_error(MycAdminHandlerErrorKind::Internal))?, + ) + .map_err(|_| handler_error(MycAdminHandlerErrorKind::Internal))?; + let url = self + .configuration + .normalized() + .pointer("/policy/challenges/url") + .and_then(Value::as_str) + .and_then(|value| MycAuthorizationChallengeUrl::new(value).ok()) + .ok_or_else(|| handler_error(MycAdminHandlerErrorKind::Unavailable))?; + let mut nonce = [0_u8; 32]; + SystemEntropy + .fill_bytes(&mut nonce) + .map_err(|_| handler_error(MycAdminHandlerErrorKind::Unavailable))?; + let challenge = MycAuthorizationChallengeRequest::new( + operation_id, + connection_id, + generation, + url, + MycAuthorizationChallengeNonce::from_injected_entropy(nonce), + issued_at, + expires_at, + ) + .map_err(|_| handler_error(MycAdminHandlerErrorKind::Internal))?; + if !self + .state + .metadata() + .admits_authorization_challenge_request(&challenge) + { + return Err(handler_error(MycAdminHandlerErrorKind::Conflict)); + } + let rate = self + .state + .metadata() + .governance_rate_policy(MycRateLimitClass::ChallengeCreation); + let completed_at = admin_operation_time(issued_at)?; + self.state + .repository() + .execute_database_admin_operation( + &request, + completed_at, + MycAdminOperationJournalPolicy::seven_days(), + move |transaction| { + Box::pin(async move { + let record = + apply_admin_challenge_require(transaction, challenge, rate).await?; + response( + route, + json!({ + "challenge_id": hex::encode(record.id().as_bytes()), + "challenge_url": record.url().as_str(), + "connection_id": hex::encode(record.connection_id().as_bytes()), + "expires_at_utc": record.expires_at().get() / 1_000, + "issued_at_utc": record.issued_at().get() / 1_000, + "request_identity": hex::encode(record.operation_id().as_bytes()), + "state": challenge_admin_state(record.state()), + }), + ) + .map_err(|_| AdminJournalOperationError::Binding) + }) + }, + ) + .await + .map_err(map_admin_journal_error) + } + + async fn challenge_authorize( + &self, + request: MycAdminRequestDocument, + ) -> Result<MycAdminResponseDocument, MycAdminHandlerError> { + let route = request.route(); + let model = request_model(&request)?; + let challenge_id = digest_id_parameter(&request, "challenge_id") + .map(MycAuthorizationChallengeId::from_bytes)?; + let generation = model_u64(&model, "/expected_generation") + .and_then(|value| MycConnectionPolicyGeneration::new(value).ok()) + .ok_or_else(|| handler_error(MycAdminHandlerErrorKind::Internal))?; + let observed_at = connection_time_now_for_admin()?; + let completed_at = admin_operation_time(observed_at)?; + let operation_id = request_operation_id(&request)?.to_owned(); + let rate = self + .state + .metadata() + .governance_rate_policy(MycRateLimitClass::ChallengeAuthorization); + self.state + .repository() + .execute_database_admin_operation( + &request, + completed_at, + MycAdminOperationJournalPolicy::seven_days(), + move |transaction| { + Box::pin(async move { + let record = apply_admin_challenge_authorize( + transaction, + challenge_id, + generation, + observed_at, + rate, + ) + .await?; + response( + route, + json!({ + "challenge_id": hex::encode(record.id().as_bytes()), + "generation": record.policy_generation().get(), + "operation_id": operation_id, + "request_identity": hex::encode(record.operation_id().as_bytes()), + "state": challenge_admin_state(record.state()), + }), + ) + .map_err(|_| AdminJournalOperationError::Binding) + }) + }, + ) + .await + .map_err(map_admin_journal_error) + } + + async fn discovery_mutation( + &self, + request: MycAdminRequestDocument, + ) -> Result<MycAdminResponseDocument, MycAdminHandlerError> { + let route = request.route(); + let model = request_model(&request)?; + let expected_generation = model_u64(&model, "/expected_generation") + .ok_or_else(|| handler_error(MycAdminHandlerErrorKind::Internal))?; + let actual_generation = u64::from( + self.state + .repository() + .current_configuration_generation() + .await + .map_err(|_| handler_error(MycAdminHandlerErrorKind::Internal))?, + ); + if expected_generation != actual_generation { + return Err(handler_error(MycAdminHandlerErrorKind::Conflict)); + } + let now_seconds = SystemWallClock + .now_utc() + .map_err(|_| handler_error(MycAdminHandlerErrorKind::Unavailable))? + .get(); + let now_millis = now_seconds + .checked_mul(1_000) + .ok_or_else(|| handler_error(MycAdminHandlerErrorKind::Internal))?; + let prepared_at = MycAdminOperationTimeUnixMs::new(now_millis) + .map_err(|_| handler_error(MycAdminHandlerErrorKind::Internal))?; + let prepared = match self + .state + .repository() + .prepare_admin_operation(&request, prepared_at) + .await + .map_err(map_admin_journal_error)? + { + MycAdminOperationAdmission::ExactReplay(response) => return Ok(response), + MycAdminOperationAdmission::Prepared(prepared) => prepared, + }; + let commit = self + .prepare_discovery_commit(&request, now_seconds, now_millis) + .await?; + let completed_at = admin_operation_time(connection_time_now_for_admin()?)?; + let operation_id = request_operation_id(&request)?.to_owned(); + match route { + MycAdminRoute::DiscoveryRender => { + let response = response( + route, + json!({ + "document_digests": discovery_digests(&commit), + "generation": actual_generation, + "operation_id": operation_id, + }), + )?; + self.state + .repository() + .complete_admin_operation( + &prepared, + &response, + completed_at, + MycAdminOperationJournalPolicy::seven_days(), + ) + .await + .map_err(map_admin_journal_error)?; + Ok(response) + } + MycAdminRoute::DiscoveryRefresh => { + let state = self + .state + .repository() + .read_discovery_publication_state() + .await + .map_err(|_| handler_error(MycAdminHandlerErrorKind::Internal))?; + let diff = if state + .as_ref() + .is_some_and(|state| state.desired_generation_id() == commit.generation_id()) + { + "in_sync" + } else { + "drift" + }; + let response = response( + route, + json!({ + "diff_state": diff, + "generation": actual_generation, + "operation_id": operation_id, + "source_completion": "complete", + }), + )?; + self.state + .repository() + .complete_admin_operation( + &prepared, + &response, + completed_at, + MycAdminOperationJournalPolicy::seven_days(), + ) + .await + .map_err(map_admin_journal_error)?; + Ok(response) + } + MycAdminRoute::DiscoveryPublish => { + let delivery_policy = self.state.metadata().delivery_policies().clone(); + let discovery_policy = self + .state + .metadata() + .discovery_policies() + .cloned() + .ok_or_else(|| handler_error(MycAdminHandlerErrorKind::Unavailable))?; + self.state + .repository() + .complete_prepared_database_admin_operation( + &prepared, + completed_at, + MycAdminOperationJournalPolicy::seven_days(), + move |transaction| { + Box::pin(async move { + let admission = apply_admin_discovery_publish( + transaction, + &commit, + &delivery_policy, + &discovery_policy, + ) + .await?; + let record = admission.record(); + response( + route, + json!({ + "exact_bytes_digest": hex::encode(commit.event_digest().as_bytes()), + "generation": actual_generation, + "operation_id": operation_id, + "target_count": u64::try_from(record.job().targets().len()).map_err(|_| AdminJournalOperationError::Binding)?, + "workflow_id": hex::encode(record.job().id().as_bytes()), + }), + ) + .map_err(|_| AdminJournalOperationError::Binding) + }) + }, + ) + .await + .map_err(map_admin_journal_error) + } + _ => Err(handler_error(MycAdminHandlerErrorKind::Internal)), + } + } + + async fn prepare_discovery_commit( + &self, + request: &MycAdminRequestDocument, + now_seconds: u64, + now_millis: u64, + ) -> Result<MycDiscoveryCommitRequest, MycAdminHandlerError> { + let unsigned = prepare_discovery_signing_bytes(self.state.metadata(), now_seconds) + .map_err(|_| handler_error(MycAdminHandlerErrorKind::Unavailable))?; + let binding = self + .configuration + .provider_contract() + .binding(MycProviderRole::Discovery) + .ok_or_else(|| handler_error(MycAdminHandlerErrorKind::Unavailable))?; + let timeout = configuration_integer( + &self.configuration, + "/transport/ingress/subscription_deadline_ms", + ) + .map_err(|_| handler_error(MycAdminHandlerErrorKind::Internal))?; + let deadline = MycProviderDeadlineUnixMs::new( + now_millis + .checked_add(timeout) + .ok_or_else(|| handler_error(MycAdminHandlerErrorKind::Internal))?, + ) + .map_err(|_| handler_error(MycAdminHandlerErrorKind::Internal))?; + let operation = MycProviderOperation::new( + binding, + MycProviderOperationId::from_bytes(admin_identity( + b"radroots.myc.admin.discovery.operation.v1\0", + request_operation_id(request)?, + )), + MycProviderCorrelationId::from_bytes(admin_identity( + b"radroots.myc.admin.discovery.correlation.v1\0", + request.correlation_id(), + )), + deadline, + MycProviderOperationInput::sign_event(&unsigned) + .map_err(|_| handler_error(MycAdminHandlerErrorKind::Internal))?, + ) + .map_err(|_| handler_error(MycAdminHandlerErrorKind::Internal))?; + let verified = self + .providers + .execute( + operation, + MycProviderResponseObservedAtUnixMs::new(now_millis) + .map_err(|_| handler_error(MycAdminHandlerErrorKind::Internal))?, + &MycTaskCancellation::uncancelled(), + ) + .await + .map_err(|_| handler_error(MycAdminHandlerErrorKind::Unavailable))?; + let event = verified + .signed_event_bytes() + .ok_or_else(|| handler_error(MycAdminHandlerErrorKind::Internal))?; + MycDiscoveryCommitRequest::new( + self.state.metadata(), + event, + MycDeliveryTimeUnixMs::new(now_millis) + .map_err(|_| handler_error(MycAdminHandlerErrorKind::Internal))?, + ) + .map_err(|_| handler_error(MycAdminHandlerErrorKind::Internal)) + } + + async fn discovery_desired( + &self, + request: MycAdminRequestDocument, + ) -> Result<MycAdminResponseDocument, MycAdminHandlerError> { + let generation = self + .state + .repository() + .current_configuration_generation() + .await + .map_err(|_| handler_error(MycAdminHandlerErrorKind::Internal))?; + let enabled = self.state.metadata().discovery_policies().is_some(); + let publication = self + .state + .repository() + .read_discovery_publication_state() + .await + .map_err(|_| handler_error(MycAdminHandlerErrorKind::Internal))?; + let document = match publication.as_ref() { + Some(publication) => self + .state + .repository() + .read_discovery_document_for_job(publication.desired_job_id()) + .await + .map_err(|_| handler_error(MycAdminHandlerErrorKind::Internal))?, + None => None, + }; + let mut digests = serde_json::Map::new(); + if let Some(document) = document { + digests.insert( + "handler_event".to_owned(), + Value::String(hex::encode(document.event_digest().as_bytes())), + ); + digests.insert( + "nip05_projection".to_owned(), + Value::String(hex::encode(document.nip05_projection_digest().as_bytes())), + ); + } + let publication_state = if !enabled { + "disabled" + } else if publication.as_ref().is_some_and(|state| { + state.current_generation_id() == Some(state.desired_generation_id()) + }) { + "delivered" + } else { + "pending" + }; + response( + request.route(), + json!({ + "document_digests": digests, + "enabled": enabled, + "generation": u64::from(generation), + "publication_state": publication_state, + }), + ) + } +} + +impl MycAdminHandler for RuntimeAdminHandler { + fn handle<'a>(&'a self, request: MycAdminRequestDocument) -> MycAdminFuture<'a> { + let mutation = request.route().is_mutation(); + Box::pin(async move { + let result = self.handle_inner(request).await; + if mutation && result.is_ok() { + let _ = self.health.try_send(HealthEvent::StateChanged); + } + result + }) + } +} + +struct ConnectionCursor { + snapshot: MycConnectionTimeUnixMs, + before: MycConnectionTimeUnixMs, + id: MycConnectionId, +} + +struct AuditCursor { + snapshot: u64, + before: u64, +} + +fn request_model(request: &MycAdminRequestDocument) -> Result<Value, MycAdminHandlerError> { + serde_json::from_slice(request.model_bytes()) + .map_err(|_| handler_error(MycAdminHandlerErrorKind::Internal)) +} + +fn request_operation_id(request: &MycAdminRequestDocument) -> Result<&str, MycAdminHandlerError> { + request + .operation_id() + .ok_or_else(|| handler_error(MycAdminHandlerErrorKind::Internal)) +} + +fn model_u64(model: &Value, pointer: &str) -> Option<u64> { + model.pointer(pointer).and_then(Value::as_u64) +} + +fn digest_id(value: &str) -> Result<[u8; 32], MycAdminHandlerError> { + let mut bytes = [0_u8; 32]; + if value.len() != 64 || hex::decode_to_slice(value, &mut bytes).is_err() { + return Err(handler_error(MycAdminHandlerErrorKind::NotFound)); + } + Ok(bytes) +} + +fn digest_id_parameter( + request: &MycAdminRequestDocument, + name: &str, +) -> Result<[u8; 32], MycAdminHandlerError> { + request + .parameter(name) + .ok_or_else(|| handler_error(MycAdminHandlerErrorKind::NotFound)) + .and_then(digest_id) +} + +fn connection_time_now_for_admin() -> Result<MycConnectionTimeUnixMs, MycAdminHandlerError> { + MycConnectionTimeUnixMs::new( + SystemWallClock + .now_utc() + .map_err(|_| handler_error(MycAdminHandlerErrorKind::Unavailable))? + .get() + .checked_mul(1_000) + .ok_or_else(|| handler_error(MycAdminHandlerErrorKind::Internal))?, + ) + .map_err(|_| handler_error(MycAdminHandlerErrorKind::Internal)) +} + +fn admin_operation_time( + time: MycConnectionTimeUnixMs, +) -> Result<MycAdminOperationTimeUnixMs, MycAdminHandlerError> { + MycAdminOperationTimeUnixMs::new(time.get()) + .map_err(|_| handler_error(MycAdminHandlerErrorKind::Internal)) +} + +fn admin_identity(domain: &[u8], value: &str) -> [u8; 32] { + let mut hasher = Sha256::new(); + hasher.update(domain); + hasher.update(u64::try_from(value.len()).unwrap_or(u64::MAX).to_be_bytes()); + hasher.update(value.as_bytes()); + hasher.finalize().into() +} + +fn admin_audit_correlation(request: &MycAdminRequestDocument) -> MycAuditCorrelationId { + MycAuditCorrelationId::new(admin_identity( + b"radroots.myc.admin.audit_correlation.v1\0", + request.correlation_id(), + )) +} + +const fn challenge_admin_state(state: MycAuthorizationChallengeState) -> &'static str { + match state { + MycAuthorizationChallengeState::Pending => "pending", + MycAuthorizationChallengeState::Authorized => "authorized", + MycAuthorizationChallengeState::Expired => "expired", + } +} + +fn discovery_digests(request: &MycDiscoveryCommitRequest) -> Value { + json!({ + "desired": hex::encode(request.desired_digest().as_bytes()), + "handler_event": hex::encode(request.event_digest().as_bytes()), + "nip05_projection": hex::encode(request.projection_digest().as_bytes()), + }) +} + +fn seconds_as_millis(seconds: u64, include_full_second: bool) -> Result<u64, MycAdminHandlerError> { + seconds + .checked_mul(1_000) + .and_then(|value| { + if include_full_second { + value.checked_add(999) + } else { + Some(value) + } + }) + .filter(|value| i64::try_from(*value).is_ok()) + .ok_or_else(|| handler_error(MycAdminHandlerErrorKind::Internal)) +} + +fn optional_seconds_as_millis( + model: &Value, + pointer: &str, + include_full_second: bool, +) -> Result<Option<u64>, MycAdminHandlerError> { + model_u64(model, pointer) + .map(|value| seconds_as_millis(value, include_full_second)) + .transpose() +} + +const fn connection_state_code(status: Option<MycConnectionStatus>) -> u8 { + match status { + None => 0, + Some(MycConnectionStatus::Pending) => 1, + Some(MycConnectionStatus::Active) => 2, + Some(MycConnectionStatus::Denied) => 3, + Some(MycConnectionStatus::Expired) => 4, + } +} + +fn encode_connection_cursor( + snapshot: MycConnectionTimeUnixMs, + before: MycConnectionTimeUnixMs, + id: MycConnectionId, + status: Option<MycConnectionStatus>, + key: &[u8; 32], +) -> String { + let mut payload = Vec::with_capacity(83); + payload.extend_from_slice(&[1, 1]); + payload.extend_from_slice(&snapshot.get().to_be_bytes()); + payload.extend_from_slice(&before.get().to_be_bytes()); + payload.extend_from_slice(id.as_bytes()); + payload.push(connection_state_code(status)); + payload.extend_from_slice(&cursor_mac(key, &payload)); + URL_SAFE_NO_PAD.encode(payload) +} + +fn decode_connection_cursor( + encoded: &str, + status: Option<MycConnectionStatus>, + key: &[u8; 32], +) -> Result<ConnectionCursor, MycAdminHandlerError> { + let bytes = URL_SAFE_NO_PAD + .decode(encoded) + .map_err(|_| handler_error(MycAdminHandlerErrorKind::InvalidCursor))?; + if bytes.len() != 83 + || bytes[0..2] != [1, 1] + || bytes[50] != connection_state_code(status) + || !constant_time_equal(&bytes[51..], &cursor_mac(key, &bytes[..51])) + { + return Err(handler_error(MycAdminHandlerErrorKind::InvalidCursor)); + } + let snapshot = MycConnectionTimeUnixMs::new(u64::from_be_bytes( + bytes[2..10] + .try_into() + .map_err(|_| handler_error(MycAdminHandlerErrorKind::InvalidCursor))?, + )) + .map_err(|_| handler_error(MycAdminHandlerErrorKind::InvalidCursor))?; + let before = MycConnectionTimeUnixMs::new(u64::from_be_bytes( + bytes[10..18] + .try_into() + .map_err(|_| handler_error(MycAdminHandlerErrorKind::InvalidCursor))?, + )) + .map_err(|_| handler_error(MycAdminHandlerErrorKind::InvalidCursor))?; + let id = MycConnectionId::from_bytes( + bytes[18..50] + .try_into() + .map_err(|_| handler_error(MycAdminHandlerErrorKind::InvalidCursor))?, + ); + Ok(ConnectionCursor { + snapshot, + before, + id, + }) +} + +fn audit_query_digest( + from: Option<u64>, + to: Option<u64>, + kind: Option<MycAuditKind>, + outcome: Option<MycAuditOutcome>, +) -> [u8; 32] { + let mut hasher = Sha256::new(); + hasher.update(b"radroots.myc.admin.audit_query.v1\0"); + hasher.update(from.unwrap_or(u64::MAX).to_be_bytes()); + hasher.update(to.unwrap_or(u64::MAX).to_be_bytes()); + hasher.update( + kind.map(MycAuditKind::as_str) + .unwrap_or_default() + .as_bytes(), + ); + hasher.update([0]); + hasher.update( + outcome + .map(MycAuditOutcome::as_str) + .unwrap_or_default() + .as_bytes(), + ); + hasher.finalize().into() +} + +fn encode_audit_cursor( + snapshot: u64, + before: u64, + query_digest: [u8; 32], + key: &[u8; 32], +) -> String { + let mut payload = Vec::with_capacity(82); + payload.extend_from_slice(&[1, 2]); + payload.extend_from_slice(&snapshot.to_be_bytes()); + payload.extend_from_slice(&before.to_be_bytes()); + payload.extend_from_slice(&query_digest); + payload.extend_from_slice(&cursor_mac(key, &payload)); + URL_SAFE_NO_PAD.encode(payload) +} + +fn decode_audit_cursor( + encoded: &str, + query_digest: [u8; 32], + key: &[u8; 32], +) -> Result<AuditCursor, MycAdminHandlerError> { + let bytes = URL_SAFE_NO_PAD + .decode(encoded) + .map_err(|_| handler_error(MycAdminHandlerErrorKind::InvalidCursor))?; + if bytes.len() != 82 + || bytes[0..2] != [1, 2] + || bytes[18..50] != query_digest + || !constant_time_equal(&bytes[50..], &cursor_mac(key, &bytes[..50])) + { + return Err(handler_error(MycAdminHandlerErrorKind::InvalidCursor)); + } + Ok(AuditCursor { + snapshot: u64::from_be_bytes( + bytes[2..10] + .try_into() + .map_err(|_| handler_error(MycAdminHandlerErrorKind::InvalidCursor))?, + ), + before: u64::from_be_bytes( + bytes[10..18] + .try_into() + .map_err(|_| handler_error(MycAdminHandlerErrorKind::InvalidCursor))?, + ), + }) +} + +fn cursor_mac(key: &[u8; 32], payload: &[u8]) -> [u8; 32] { + let mut hasher = Sha256::new(); + hasher.update(b"radroots.myc.admin.cursor.v1\0"); + hasher.update(key); + hasher.update( + u64::try_from(payload.len()) + .unwrap_or(u64::MAX) + .to_be_bytes(), + ); + hasher.update(payload); + hasher.finalize().into() +} + +fn constant_time_equal(left: &[u8], right: &[u8]) -> bool { + left.len() == right.len() + && left + .iter() + .zip(right) + .fold(0_u8, |difference, (left, right)| { + difference | (left ^ right) + }) + == 0 +} + +fn safe_metric_key(value: &str) -> bool { + (1..=64).contains(&value.len()) + && value.as_bytes().first().is_some_and(u8::is_ascii_lowercase) + && value + .bytes() + .all(|byte| byte.is_ascii_lowercase() || byte.is_ascii_digit() || byte == b'_') +} + +struct RuntimeShutdownHandler { + accepting_mutations: Arc<AtomicBool>, + state: Arc<MycStateHost>, + status: MycStatusPublisher, + status_context: RuntimeStatusContext, + health: RuntimeHealth, +} + +impl RuntimeShutdownHandler { + fn publish(&mut self, phase: MycServicePhase) -> Result<(), HostError> { + let observation = self + .status_context + .observation(phase, &self.health) + .map_err(|error| HostError::with_source(HostErrorKind::Lifecycle, error))?; + self.status + .publish(observation) + .map_err(|error| HostError::with_source(HostErrorKind::Lifecycle, error)) + } +} + +impl ShutdownPhaseHandler for RuntimeShutdownHandler { + fn enter(&mut self, phase: ShutdownPhase) -> ShutdownPhaseFuture<'_> { + Box::pin(async move { + match phase { + ShutdownPhase::RejectNewMutations => { + self.accepting_mutations.store(false, Ordering::Release); + self.publish(MycServicePhase::Stopping)?; + } + ShutdownPhase::PersistRecoverableWork => { + self.state + .repository() + .verify_delivery_invariants() + .await + .map_err(|error| HostError::with_source(HostErrorKind::Lifecycle, error))?; + } + ShutdownPhase::CloseSqlite => { + self.state + .close() + .await + .map_err(|error| HostError::with_source(HostErrorKind::Lifecycle, error))?; + } + ShutdownPhase::CancelIngress + | ShutdownPhase::DrainOperations + | ShutdownPhase::CloseNetwork + | ShutdownPhase::CloseSockets => {} + } + Ok(()) + }) + } +} + +pub(crate) async fn run_myc_daemon<S>( + runtime: MycRuntimeContext, + configuration: MycConfigDocumentV1, + applied_at: MigrationAppliedAtUnixSeconds, + build: &MigrationBuildIdentity, + signals: S, +) -> MycProcessResult +where + S: MycProcessSignalSource + 'static, +{ + run_myc_daemon_inner(runtime, configuration, applied_at, build, signals) + .await + .unwrap_or_else(|error| error.process_result()) +} + +async fn run_myc_daemon_inner<S>( + runtime: MycRuntimeContext, + configuration: MycConfigDocumentV1, + applied_at: MigrationAppliedAtUnixSeconds, + build: &MigrationBuildIdentity, + signals: S, +) -> Result<MycProcessResult, MycDaemonError> +where + S: MycProcessSignalSource + 'static, +{ + let configuration = Arc::new(configuration); + let state = Arc::new( + open_myc_state_read_write_from_config(&runtime, &configuration, applied_at, build) + .await + .map_err(|_| MycDaemonError::new(MycDaemonErrorKind::State))?, + ); + let now_millis = wall_time_millis()?; + let mut recovery_entropy = [0_u8; 32]; + SystemEntropy + .fill_bytes(&mut recovery_entropy) + .map_err(|_| MycDaemonError::new(MycDaemonErrorKind::Runtime))?; + state + .repository() + .recover_delivery_state( + MycDeliveryTimeUnixMs::new(now_millis) + .map_err(|_| MycDaemonError::new(MycDaemonErrorKind::Runtime))?, + MycDeliveryRecoveryEntropy::from_injected_entropy(recovery_entropy), + ) + .await + .map_err(|_| MycDaemonError::new(MycDaemonErrorKind::State))?; + state + .repository() + .verify_delivery_invariants() + .await + .map_err(|_| MycDaemonError::new(MycDaemonErrorKind::State))?; + + let startup_cancellation = MycTaskCancellation::uncancelled(); + let providers = Arc::new( + MycProviderExecutor::open(&runtime, &configuration, &startup_cancellation) + .await + .map_err(|_| MycDaemonError::new(MycDaemonErrorKind::Provider))?, + ); + let mut provider_seed = [0_u8; 32]; + SystemEntropy + .fill_bytes(&mut provider_seed) + .map_err(|_| MycDaemonError::new(MycDaemonErrorKind::Runtime))?; + providers + .probe_all(now_millis, provider_seed, &startup_cancellation) + .await + .map_err(|_| MycDaemonError::new(MycDaemonErrorKind::Provider))?; + let ingress_adapter = Arc::new( + MycNostrIngressAdapter::from_configuration(&configuration) + .map_err(|_| MycDaemonError::new(MycDaemonErrorKind::Relay))?, + ); + let ingress_slots = open_initial_subscriptions(&ingress_adapter, &configuration).await?; + + let generation = u64::from( + state + .repository() + .current_configuration_generation() + .await + .map_err(|_| MycDaemonError::new(MycDaemonErrorKind::State))?, + ); + let mut status_context = RuntimeStatusContext { + runtime: runtime.clone(), + configuration: Arc::clone(&configuration), + state: Arc::clone(&state), + clock: SystemMonotonicClock::new(), + generation, + connection_counts: MycConnectionCountsV1::default(), + outbox: MycOutboxStatusV1::default(), + }; + status_context.refresh_state().await?; + let operations_enabled = configuration + .normalized() + .pointer("/operations/enabled") + .and_then(Value::as_bool) + .unwrap_or(false); + let health = RuntimeHealth::ready(true); + let (publisher, status_reader) = myc_status_cache( + runtime.context().instance().clone(), + status_context.observation(MycServicePhase::Starting, &health)?, + ) + .map_err(|_| MycDaemonError::new(MycDaemonErrorKind::Runtime))?; + let accepting_mutations = Arc::new(AtomicBool::new(true)); + let mut cursor_key = [0_u8; 32]; + SystemEntropy + .fill_bytes(&mut cursor_key) + .map_err(|_| MycDaemonError::new(MycDaemonErrorKind::Runtime))?; + let (health_tx, mut health_rx) = mpsc::channel(8); + let admin_handler = Arc::new(RuntimeAdminHandler { + state: Arc::clone(&state), + configuration: Arc::clone(&configuration), + providers: Arc::clone(&providers), + status: status_reader.clone(), + accepting_mutations: Arc::clone(&accepting_mutations), + cursor_key, + health: health_tx.clone(), + }); + let admin = MycAdminServer::new(&configuration, admin_handler) + .map_err(|_| MycDaemonError::new(MycDaemonErrorKind::Admin))? + .bind(&runtime) + .await + .map_err(|_| MycDaemonError::new(MycDaemonErrorKind::Admin))?; + let operations = if operations_enabled { + Some( + MycOperationsServer::new(&configuration, &status_reader) + .map_err(|_| MycDaemonError::new(MycDaemonErrorKind::Operations))? + .bind() + .await + .map_err(|_| MycDaemonError::new(MycDaemonErrorKind::Operations))?, + ) + } else { + None + }; + + let mut shutdown_handler = RuntimeShutdownHandler { + accepting_mutations, + state: Arc::clone(&state), + status: publisher, + status_context, + health, + }; + let ingress_capacity = configuration_usize(&configuration, "/resource_limits/queues/ingress")?; + let provider_capacity = + configuration_usize(&configuration, "/resource_limits/queues/provider")?; + let (ingress_tx, ingress_rx) = mpsc::channel(ingress_capacity); + let mut supervisor = TaskSupervisor::new(); + spawn_admin(&mut supervisor, admin)?; + if let Some(operations) = operations { + spawn_operations(&mut supervisor, operations)?; + } + spawn_relay_ingress( + &mut supervisor, + ingress_adapter, + ingress_slots, + Arc::clone(&configuration), + health_tx.clone(), + ingress_tx, + )?; + spawn_provider_dispatch( + &mut supervisor, + Arc::clone(&providers), + Arc::clone(&configuration), + Arc::clone(&state), + health_tx.clone(), + ingress_rx, + provider_capacity, + )?; + spawn_delivery_outbox( + &mut supervisor, + Arc::clone(&state), + Arc::clone(&configuration), + health_tx, + )?; + shutdown_handler + .publish(MycServicePhase::Ready) + .map_err(|_| MycDaemonError::new(MycDaemonErrorKind::Runtime))?; + + let mut signals = ProcessSignalAdapter::new(HostSignalSource::new(signals)); + let signal_initiated = loop { + tokio::select! { + action = signals.next_action() => { + match action { + Ok(action) if !action.forces_termination() => break true, + Ok(_) | Err(_) => break false, + } + } + joined = supervisor.join_next() => { + match joined { + Some(Ok(exit)) if exit.status() == SupervisedTaskExitStatus::OptionalFailure + || exit.metadata().name().as_str() == TASK_OPERATIONS_SERVER => { + shutdown_handler.health.operations = false; + shutdown_handler.publish(MycServicePhase::Degraded) + .map_err(|_| MycDaemonError::new(MycDaemonErrorKind::Runtime))?; + } + Some(Ok(_)) | Some(Err(_)) | None => break false, + } + } + event = health_rx.recv() => { + let Some(event) = event else { break false; }; + if shutdown_handler.health.observe(event) { + if event == HealthEvent::StateChanged { + shutdown_handler.status_context.refresh_state().await?; + } + let phase = if shutdown_handler.health.providers + && shutdown_handler.health.required_relays { + if shutdown_handler.health.operations { + MycServicePhase::Ready + } else { + MycServicePhase::Degraded + } + } else { + MycServicePhase::Unready + }; + shutdown_handler.publish(phase) + .map_err(|_| MycDaemonError::new(MycDaemonErrorKind::Runtime))?; + } + } + } + }; + let grace = Duration::from_millis(configuration_integer( + &configuration, + "/service/shutdown_grace_ms", + )?); + let mut shutdown = GracefulShutdown::new(grace) + .map_err(|_| MycDaemonError::new(MycDaemonErrorKind::Runtime))?; + let clock = SystemMonotonicClock::new(); + let summary = if signal_initiated { + shutdown + .run(&clock, &mut supervisor, &mut shutdown_handler, async { + let _ = signals.next_action().await; + }) + .await + } else { + shutdown + .run( + &clock, + &mut supervisor, + &mut shutdown_handler, + pending::<()>(), + ) + .await + } + .map_err(|_| MycDaemonError::new(MycDaemonErrorKind::Runtime))?; + if signal_initiated && summary.disposition() == ShutdownDisposition::Completed { + Ok(MycProcessResult::Success) + } else { + Err(MycDaemonError::new(MycDaemonErrorKind::Runtime)) + } +} + +fn spawn_admin( + supervisor: &mut TaskSupervisor, + server: MycBoundAdminServer, +) -> Result<(), MycDaemonError> { + supervisor + .spawn(task_metadata(TASK_ADMIN_SERVER, TaskClassification::Critical, ShutdownPhase::CloseSockets)?, move |cancellation| async move { + let token = MycAdminCancellationToken::new(); + let serve_token = token.clone(); + let serve = server.serve(serve_token); + tokio::pin!(serve); + tokio::select! { + result = serve.as_mut() => result.map_err(|error| HostError::with_source(HostErrorKind::AdminTransport, error)), + () = cancellation.cancelled() => { + token.cancel(); + serve.await.map_err(|error| HostError::with_source(HostErrorKind::AdminTransport, error)) + } + } + }) + .map_err(|_| MycDaemonError::new(MycDaemonErrorKind::Runtime)) +} + +fn spawn_operations( + supervisor: &mut TaskSupervisor, + server: MycBoundOperationsServer, +) -> Result<(), MycDaemonError> { + supervisor + .spawn(task_metadata(TASK_OPERATIONS_SERVER, TaskClassification::Optional, ShutdownPhase::CloseSockets)?, move |cancellation| async move { + let token = MycOperationsCancellationToken::new(); + let serve_token = token.clone(); + let serve = server.serve(serve_token); + tokio::pin!(serve); + tokio::select! { + result = serve.as_mut() => result.map_err(|error| HostError::with_source(HostErrorKind::OperationsServe, error)), + () = cancellation.cancelled() => { + token.cancel(); + serve.await.map_err(|error| HostError::with_source(HostErrorKind::OperationsServe, error)) + } + } + }) + .map_err(|_| MycDaemonError::new(MycDaemonErrorKind::Runtime)) +} + +fn spawn_relay_ingress( + supervisor: &mut TaskSupervisor, + adapter: Arc<MycNostrIngressAdapter>, + slots: Vec<RuntimeIngressSlot>, + configuration: Arc<MycConfigDocumentV1>, + health: mpsc::Sender<HealthEvent>, + ingress: mpsc::Sender<RuntimeIngressItem>, +) -> Result<(), MycDaemonError> { + supervisor + .spawn( + task_metadata( + TASK_RELAY_INGRESS, + TaskClassification::Critical, + ShutdownPhase::CancelIngress, + )?, + move |cancellation| async move { + run_relay_ingress(cancellation, adapter, slots, configuration, health, ingress) + .await + }, + ) + .map_err(|_| MycDaemonError::new(MycDaemonErrorKind::Runtime)) +} + +fn spawn_provider_dispatch( + supervisor: &mut TaskSupervisor, + providers: Arc<MycProviderExecutor>, + configuration: Arc<MycConfigDocumentV1>, + state: Arc<MycStateHost>, + health: mpsc::Sender<HealthEvent>, + ingress: mpsc::Receiver<RuntimeIngressItem>, + provider_capacity: usize, +) -> Result<(), MycDaemonError> { + supervisor + .spawn( + task_metadata( + TASK_PROVIDER_DISPATCH, + TaskClassification::Critical, + ShutdownPhase::DrainOperations, + )?, + move |cancellation| async move { + run_provider_dispatch( + cancellation, + configuration, + state, + providers, + health, + ingress, + provider_capacity, + ) + .await + }, + ) + .map_err(|_| MycDaemonError::new(MycDaemonErrorKind::Runtime)) +} + +async fn open_initial_subscriptions( + adapter: &Arc<MycNostrIngressAdapter>, + configuration: &MycConfigDocumentV1, +) -> Result<Vec<RuntimeIngressSlot>, MycDaemonError> { + let initial = + configuration_integer(configuration, "/transport/publish_retry/initial_backoff_ms")?; + let mut slots = Vec::with_capacity(adapter.group_count()); + for index in 0..adapter.group_count() { + let required = adapter + .group_is_required(index) + .ok_or_else(|| MycDaemonError::new(MycDaemonErrorKind::Relay))?; + let duration = + configuration_integer(configuration, "/transport/ingress/subscription_deadline_ms")?; + let result = subscribe_ingress(Arc::clone(adapter), index, Vec::new(), duration).await; + let state = match result { + Ok(subscription) => RuntimeIngressSlotState::Active(subscription), + Err(_) if required => return Err(MycDaemonError::new(MycDaemonErrorKind::Relay)), + Err(_) => RuntimeIngressSlotState::Waiting, + }; + slots.push(RuntimeIngressSlot { + index, + required, + checkpoints: Vec::new(), + state, + retry_delay_ms: initial, + }); + } + if slots.is_empty() + || slots + .iter() + .any(|slot| slot.required && !matches!(slot.state, RuntimeIngressSlotState::Active(_))) + { + return Err(MycDaemonError::new(MycDaemonErrorKind::Relay)); + } + Ok(slots) +} + +async fn subscribe_ingress( + adapter: Arc<MycNostrIngressAdapter>, + index: usize, + checkpoints: Vec<SubscriptionCheckpoint>, + subscription_duration_ms: u64, +) -> Result<BoxSubscription, crate::transport_nostr_adapter::MycRelayAdapterError> { + let deadline = wall_time_millis() + .ok() + .and_then(|now| now.checked_add(subscription_duration_ms)) + .ok_or_else(|| { + crate::transport_nostr_adapter::runtime_relay_adapter_error( + crate::transport_nostr_adapter::MycRelayAdapterErrorKind::Configuration, + ) + })?; + adapter + .subscribe(index, 1_000, deadline, &checkpoints) + .await +} + +async fn run_relay_ingress( + cancellation: radroots_service_host::CancellationToken, + adapter: Arc<MycNostrIngressAdapter>, + mut slots: Vec<RuntimeIngressSlot>, + configuration: Arc<MycConfigDocumentV1>, + health: mpsc::Sender<HealthEvent>, + ingress: mpsc::Sender<RuntimeIngressItem>, +) -> Result<(), HostError> { + if slots.is_empty() || slots.len() > 2 { + return Err(HostError::new(HostErrorKind::TaskFailure)); + } + loop { + if cancellation.is_cancelled() { + cancel_ingress_slots(&mut slots).await; + return Ok(()); + } + let selected = if slots.len() == 1 { + tokio::select! { + () = cancellation.cancelled() => { + cancel_ingress_slots(&mut slots).await; + return Ok(()); + } + action = poll_ingress_slot(&mut slots[0]) => (0, action), + } + } else { + let (left, right) = slots.split_at_mut(1); + tokio::select! { + () = cancellation.cancelled() => { + cancel_ingress_slots(&mut slots).await; + return Ok(()); + } + action = poll_ingress_slot(&mut left[0]) => (0, action), + action = poll_ingress_slot(&mut right[0]) => (1, action), + } + }; + handle_ingress_action( + selected.0, + selected.1, + &adapter, + &configuration, + &mut slots, + &health, + &ingress, + &cancellation, + ) + .await?; + } +} + +async fn poll_ingress_slot(slot: &mut RuntimeIngressSlot) -> RuntimeIngressAction { + match &mut slot.state { + RuntimeIngressSlotState::Active(subscription) => { + RuntimeIngressAction::Next(subscription.next().await) + } + RuntimeIngressSlotState::Waiting => { + tokio::time::sleep(Duration::from_millis(slot.retry_delay_ms)).await; + RuntimeIngressAction::Retry + } + RuntimeIngressSlotState::Connecting(connecting) => { + RuntimeIngressAction::Connected(connecting.await) + } + } +} + +#[allow(clippy::too_many_arguments)] +async fn handle_ingress_action( + selected: usize, + action: RuntimeIngressAction, + adapter: &Arc<MycNostrIngressAdapter>, + configuration: &MycConfigDocumentV1, + slots: &mut [RuntimeIngressSlot], + health: &mpsc::Sender<HealthEvent>, + ingress: &mpsc::Sender<RuntimeIngressItem>, + cancellation: &radroots_service_host::CancellationToken, +) -> Result<(), HostError> { + let initial = + configuration_integer(configuration, "/transport/publish_retry/initial_backoff_ms") + .map_err(|error| HostError::with_source(HostErrorKind::TaskFailure, error))?; + let maximum = + configuration_integer(configuration, "/transport/publish_retry/maximum_backoff_ms") + .map_err(|error| HostError::with_source(HostErrorKind::TaskFailure, error))?; + match action { + RuntimeIngressAction::Next(Ok(SubscriptionNext::Event(event))) => { + let observed = event.observed(); + let relay_id = adapter + .relay_id(selected, observed.provenance().target()) + .map_err(|error| HostError::with_source(HostErrorKind::TaskFailure, error))?; + replace_checkpoint(&mut slots[selected].checkpoints, event.checkpoint().clone()); + let item = RuntimeIngressItem { + raw_event: Box::from(observed.event().raw_json().as_bytes()), + relay_id, + observed_at_unix_ms: observed.provenance().observed_at_unix_ms(), + }; + tokio::select! { + result = ingress.send(item) => result.map_err(|_| HostError::new(HostErrorKind::TaskFailure))?, + () = cancellation.cancelled() => return Ok(()), + } + } + RuntimeIngressAction::Next(Ok(SubscriptionNext::End(end))) => { + slots[selected].checkpoints = end.checkpoints().to_vec(); + slots[selected].state = RuntimeIngressSlotState::Waiting; + send_required_relay_health(required_relays_ready(slots), health, cancellation).await?; + } + RuntimeIngressAction::Next(Err(_)) => { + slots[selected].state = RuntimeIngressSlotState::Waiting; + send_required_relay_health(required_relays_ready(slots), health, cancellation).await?; + } + RuntimeIngressAction::Retry => { + let adapter = Arc::clone(adapter); + let index = slots[selected].index; + let checkpoints = slots[selected].checkpoints.clone(); + let duration = + configuration_integer(configuration, "/transport/ingress/subscription_deadline_ms") + .map_err(|error| HostError::with_source(HostErrorKind::TaskFailure, error))?; + slots[selected].state = RuntimeIngressSlotState::Connecting(Box::pin(async move { + subscribe_ingress(adapter, index, checkpoints, duration).await + })); + } + RuntimeIngressAction::Connected(Ok(subscription)) => { + slots[selected].state = RuntimeIngressSlotState::Active(subscription); + slots[selected].retry_delay_ms = initial; + send_required_relay_health(required_relays_ready(slots), health, cancellation).await?; + } + RuntimeIngressAction::Connected(Err(_)) => { + slots[selected].state = RuntimeIngressSlotState::Waiting; + slots[selected].retry_delay_ms = slots[selected] + .retry_delay_ms + .saturating_mul(2) + .min(maximum); + send_required_relay_health(required_relays_ready(slots), health, cancellation).await?; + } + } + Ok(()) +} + +fn replace_checkpoint( + checkpoints: &mut Vec<SubscriptionCheckpoint>, + checkpoint: SubscriptionCheckpoint, +) { + if let Some(existing) = checkpoints + .iter_mut() + .find(|existing| existing.target() == checkpoint.target()) + { + *existing = checkpoint; + } else { + checkpoints.push(checkpoint); + } +} + +fn required_relays_ready(slots: &[RuntimeIngressSlot]) -> bool { + slots + .iter() + .all(|slot| !slot.required || matches!(slot.state, RuntimeIngressSlotState::Active(_))) +} + +async fn send_required_relay_health( + ready: bool, + health: &mpsc::Sender<HealthEvent>, + cancellation: &radroots_service_host::CancellationToken, +) -> Result<(), HostError> { + tokio::select! { + result = health.send(HealthEvent::RequiredRelays(ready)) => { + result.map_err(|_| HostError::new(HostErrorKind::TaskFailure)) + } + () = cancellation.cancelled() => Ok(()), + } +} + +async fn cancel_ingress_slots(slots: &mut [RuntimeIngressSlot]) { + for slot in slots { + if let RuntimeIngressSlotState::Active(subscription) = &mut slot.state { + let _ = subscription.cancel().await; + } + } +} + +async fn run_provider_dispatch( + cancellation: radroots_service_host::CancellationToken, + configuration: Arc<MycConfigDocumentV1>, + state: Arc<MycStateHost>, + providers: Arc<MycProviderExecutor>, + health: mpsc::Sender<HealthEvent>, + mut ingress: mpsc::Receiver<RuntimeIngressItem>, + provider_capacity: usize, +) -> Result<(), HostError> { + if provider_capacity == 0 { + return Err(HostError::new(HostErrorKind::TaskFailure)); + } + let coordinator = MycRuntimeNip46Coordinator::new(configuration.clone(), state, providers) + .map_err(|error| HostError::with_source(HostErrorKind::TaskFailure, error))?; + let task_cancellation = MycTaskCancellation::from_host(cancellation.clone()); + let initial = configuration_integer( + &configuration, + "/transport/publish_retry/initial_backoff_ms", + ) + .map_err(|error| HostError::with_source(HostErrorKind::TaskFailure, error))?; + let maximum = configuration_integer( + &configuration, + "/transport/publish_retry/maximum_backoff_ms", + ) + .map_err(|error| HostError::with_source(HostErrorKind::TaskFailure, error))?; + let mut queue = VecDeque::with_capacity(provider_capacity); + loop { + if queue.is_empty() { + tokio::select! { + item = ingress.recv() => match item { + Some(item) => queue.push_back(item), + None => return Err(HostError::new(HostErrorKind::TaskFailure)), + }, + () = cancellation.cancelled() => return Ok(()), + } + } + while queue.len() < provider_capacity { + match ingress.try_recv() { + Ok(item) => queue.push_back(item), + Err(mpsc::error::TryRecvError::Empty) => break, + Err(mpsc::error::TryRecvError::Disconnected) if queue.is_empty() => { + return Err(HostError::new(HostErrorKind::TaskFailure)); + } + Err(mpsc::error::TryRecvError::Disconnected) => break, + } + } + let item = queue.pop_front().expect("provider queue is nonempty"); + let admission_evidence = runtime_nip46_admission_evidence() + .map_err(|error| HostError::with_source(HostErrorKind::TaskFailure, error))?; + let mut retry = initial; + loop { + if cancellation.is_cancelled() { + return Ok(()); + } + match coordinator + .process( + &item.raw_event, + item.relay_id.clone(), + item.observed_at_unix_ms, + admission_evidence, + &task_cancellation, + ) + .await + { + Ok(_) => { + let _ = health.try_send(HealthEvent::StateChanged); + send_provider_health(&health, true, &cancellation).await?; + break; + } + Err(error) if error.kind() == MycNip46DispatchErrorKind::Provider => { + send_provider_health(&health, false, &cancellation).await?; + tokio::select! { + () = cancellation.cancelled() => return Ok(()), + () = tokio::time::sleep(Duration::from_millis(retry)) => {} + } + retry = retry.saturating_mul(2).min(maximum); + } + Err(error) => { + return Err(HostError::with_source(HostErrorKind::TaskFailure, error)); + } + } + } + } +} + +fn runtime_nip46_admission_evidence() -> Result<MycRuntimeNip46AdmissionEvidence, MycDaemonError> { + let mut request_nonce = [0_u8; 32]; + SystemEntropy + .fill_bytes(&mut request_nonce) + .map_err(|_| MycDaemonError::new(MycDaemonErrorKind::Runtime))?; + MycRuntimeNip46AdmissionEvidence::new(request_nonce, wall_time_millis()?) + .map_err(|_| MycDaemonError::new(MycDaemonErrorKind::Runtime)) +} + +async fn send_provider_health( + health: &mpsc::Sender<HealthEvent>, + ready: bool, + cancellation: &radroots_service_host::CancellationToken, +) -> Result<(), HostError> { + tokio::select! { + result = health.send(HealthEvent::Providers(ready)) => { + result.map_err(|_| HostError::new(HostErrorKind::TaskFailure)) + } + () = cancellation.cancelled() => Ok(()), + } +} + +fn spawn_delivery_outbox( + supervisor: &mut TaskSupervisor, + state: Arc<MycStateHost>, + configuration: Arc<MycConfigDocumentV1>, + health: mpsc::Sender<HealthEvent>, +) -> Result<(), MycDaemonError> { + supervisor + .spawn( + task_metadata( + TASK_DELIVERY_OUTBOX, + TaskClassification::Critical, + ShutdownPhase::PersistRecoverableWork, + )?, + move |cancellation| async move { + let worker = MycDeliveryWorker::from_configuration(&configuration) + .map_err(|error| HostError::with_source(HostErrorKind::TaskFailure, error))?; + let task_cancellation = MycTaskCancellation::from_host(cancellation.clone()); + let idle = Duration::from_millis( + configuration_integer( + &configuration, + "/transport/publish_retry/initial_backoff_ms", + ) + .map_err(|error| HostError::with_source(HostErrorKind::TaskFailure, error))?, + ); + loop { + if cancellation.is_cancelled() { + return Ok(()); + } + let now = wall_time_millis().map_err(|error| { + HostError::with_source(HostErrorKind::TaskFailure, error) + })?; + let now = MycDeliveryTimeUnixMs::new(now).map_err(|error| { + HostError::with_source(HostErrorKind::TaskFailure, error) + })?; + let next = state + .repository() + .next_ready_delivery_target(now) + .await + .map_err(|error| { + HostError::with_source(HostErrorKind::TaskFailure, error) + })?; + let Some((job, relay)) = next else { + tokio::select! { + () = cancellation.cancelled() => return Ok(()), + () = tokio::time::sleep(idle) => continue, + } + }; + let mut nonce = [0_u8; 32]; + SystemEntropy.fill_bytes(&mut nonce).map_err(|error| { + HostError::with_source(HostErrorKind::TaskFailure, error) + })?; + let mut retry_entropy = [0_u8; 8]; + SystemEntropy + .fill_bytes(&mut retry_entropy) + .map_err(|error| { + HostError::with_source(HostErrorKind::TaskFailure, error) + })?; + let evidence = MycDeliveryExecutionEvidence { + claimed_at: now, + submitted_at: now, + observed_at: now, + retry_entropy, + }; + worker + .run_one( + &state.repository(), + job, + &relay, + MycDeliveryAttemptNonce::from_injected_entropy(nonce), + evidence, + &task_cancellation, + ) + .await + .map_err(|error| { + HostError::with_source(HostErrorKind::TaskFailure, error) + })?; + let _ = health.try_send(HealthEvent::StateChanged); + } + }, + ) + .map_err(|_| MycDaemonError::new(MycDaemonErrorKind::Runtime)) +} + +fn task_metadata( + name: &'static str, + classification: TaskClassification, + phase: ShutdownPhase, +) -> Result<TaskMetadata, MycDaemonError> { + TaskMetadata::new( + TaskName::new(name).map_err(|_| MycDaemonError::new(MycDaemonErrorKind::Runtime))?, + classification, + Some(phase), + ) + .map_err(|_| MycDaemonError::new(MycDaemonErrorKind::Runtime)) +} + +fn runtime_build_info() -> Result<MycStatusBuildInfoV1, MycDaemonError> { + let service_revision = option_env!("RADROOTS_SERVICE_REVISION"); + let lib_revision = option_env!("RADROOTS_LIB_REVISION"); + let rust_version = option_env!("RADROOTS_RUST_VERSION"); + let target = option_env!("RADROOTS_BUILD_TARGET"); + let mode = if service_revision.is_some() + && lib_revision.is_some() + && rust_version.is_some() + && target.is_some() + { + MycStatusBuildMode::Release + } else { + MycStatusBuildMode::Development + }; + MycStatusBuildInfoV1::new( + mode, + Some(env!("CARGO_PKG_VERSION")), + service_revision, + lib_revision, + rust_version, + target, + Some("service-host"), + ) + .map_err(|_| MycDaemonError::new(MycDaemonErrorKind::Runtime)) +} + +fn wall_time_millis() -> Result<u64, MycDaemonError> { + SystemWallClock + .now_utc() + .map_err(|_| MycDaemonError::new(MycDaemonErrorKind::Runtime))? + .get() + .checked_mul(1_000) + .filter(|value| *value <= i64::MAX as u64) + .ok_or_else(|| MycDaemonError::new(MycDaemonErrorKind::Runtime)) +} + +fn configuration_integer( + configuration: &MycConfigDocumentV1, + pointer: &str, +) -> Result<u64, MycDaemonError> { + configuration + .normalized() + .pointer(pointer) + .and_then(Value::as_u64) + .ok_or_else(|| MycDaemonError::new(MycDaemonErrorKind::Runtime)) +} + +fn configuration_usize( + configuration: &MycConfigDocumentV1, + pointer: &str, +) -> Result<usize, MycDaemonError> { + configuration_integer(configuration, pointer)? + .try_into() + .map_err(|_| MycDaemonError::new(MycDaemonErrorKind::Runtime)) +} + +fn response( + route: MycAdminRoute, + value: Value, +) -> Result<MycAdminResponseDocument, MycAdminHandlerError> { + let bytes = serde_json::to_vec(&value) + .map_err(|_| handler_error(MycAdminHandlerErrorKind::Internal))?; + MycAdminResponseDocument::from_canonical_bytes(route, &bytes) + .map_err(|_| handler_error(MycAdminHandlerErrorKind::Internal)) +} + +const fn handler_error(kind: MycAdminHandlerErrorKind) -> MycAdminHandlerError { + MycAdminHandlerError::new(kind) +} + +const fn map_admin_journal_error(error: crate::MycAdminOperationError) -> MycAdminHandlerError { + match error.kind() { + MycAdminOperationErrorKind::OperationConflict => { + handler_error(MycAdminHandlerErrorKind::OperationIdConflict) + } + MycAdminOperationErrorKind::OperationOutcomeUnknown + | MycAdminOperationErrorKind::CommitOutcomeUnknown => { + handler_error(MycAdminHandlerErrorKind::Conflict) + } + MycAdminOperationErrorKind::ResourceExhausted => { + handler_error(MycAdminHandlerErrorKind::Unavailable) + } + MycAdminOperationErrorKind::InvalidMode + | MycAdminOperationErrorKind::InvalidInput + | MycAdminOperationErrorKind::Binding + | MycAdminOperationErrorKind::Transaction => { + handler_error(MycAdminHandlerErrorKind::Internal) + } + } +} + +#[cfg(test)] +mod tests { + use std::{fs, os::unix::fs::PermissionsExt as _}; + + use super::*; + + #[test] + fn exact_task_inventory_and_shutdown_phases_are_frozen() { + let expected = [ + ( + TASK_ADMIN_SERVER, + TaskClassification::Critical, + ShutdownPhase::CloseSockets, + ), + ( + TASK_OPERATIONS_SERVER, + TaskClassification::Optional, + ShutdownPhase::CloseSockets, + ), + ( + TASK_RELAY_INGRESS, + TaskClassification::Critical, + ShutdownPhase::CancelIngress, + ), + ( + TASK_PROVIDER_DISPATCH, + TaskClassification::Critical, + ShutdownPhase::DrainOperations, + ), + ( + TASK_DELIVERY_OUTBOX, + TaskClassification::Critical, + ShutdownPhase::PersistRecoverableWork, + ), + ]; + for (name, classification, phase) in expected { + let metadata = task_metadata(name, classification, phase).expect("task metadata"); + assert_eq!(metadata.name().as_str(), name); + assert_eq!(metadata.classification(), classification); + assert_eq!(metadata.shutdown_phase(), Some(phase)); + } + } + + #[test] + fn runtime_errors_are_safe_and_have_stable_exit_classes() { + for (kind, result) in [ + ( + MycDaemonErrorKind::State, + MycProcessResult::StateOrIdentityUnavailable, + ), + ( + MycDaemonErrorKind::Provider, + MycProcessResult::ServiceOrDependencyUnavailable, + ), + ( + MycDaemonErrorKind::Relay, + MycProcessResult::ServiceOrDependencyUnavailable, + ), + ( + MycDaemonErrorKind::Admin, + MycProcessResult::ServiceOrDependencyUnavailable, + ), + ( + MycDaemonErrorKind::Operations, + MycProcessResult::ServiceOrDependencyUnavailable, + ), + ( + MycDaemonErrorKind::Runtime, + MycProcessResult::UnexpectedInternal, + ), + ] { + let error = MycDaemonError::new(kind); + assert_eq!(error.process_result(), result); + assert_eq!(error.to_string(), "Myc daemon failed"); + assert!(Error::source(&error).is_none()); + } + } + + #[test] + fn backoff_inputs_are_bounded_by_the_admitted_configuration() { + let configuration = crate::parse_myc_config_v1( + include_bytes!("../contracts/services_hardening/config.v1.example.toml"), + crate::MycConfigProfile::Production, + ) + .expect("configuration"); + assert_eq!( + configuration_integer( + &configuration, + "/transport/publish_retry/initial_backoff_ms" + ) + .expect("initial"), + 250 + ); + assert_eq!( + configuration_integer( + &configuration, + "/transport/publish_retry/maximum_backoff_ms" + ) + .expect("maximum"), + 30_000 + ); + } + + #[test] + fn admin_cursors_are_authenticated_query_bound_and_round_trip_exactly() { + let key = [0x42; 32]; + let snapshot = MycConnectionTimeUnixMs::new(10_000).expect("snapshot"); + let before = MycConnectionTimeUnixMs::new(9_000).expect("before"); + let id = MycConnectionId::from_bytes([0x24; 32]); + let encoded = encode_connection_cursor( + snapshot, + before, + id, + Some(MycConnectionStatus::Active), + &key, + ); + let decoded = decode_connection_cursor(&encoded, Some(MycConnectionStatus::Active), &key) + .expect("connection cursor"); + assert_eq!(decoded.snapshot, snapshot); + assert_eq!(decoded.before, before); + assert_eq!(decoded.id, id); + assert!( + decode_connection_cursor(&encoded, Some(MycConnectionStatus::Denied), &key).is_err() + ); + + let query = audit_query_digest( + Some(1_000), + Some(9_999), + Some(MycAuditKind::ChallengeAuthorization), + Some(MycAuditOutcome::Succeeded), + ); + let encoded = encode_audit_cursor(17, 11, query, &key); + let decoded = decode_audit_cursor(&encoded, query, &key).expect("audit cursor"); + assert_eq!(decoded.snapshot, 17); + assert_eq!(decoded.before, 11); + let other_query = audit_query_digest(None, None, None, None); + assert!(decode_audit_cursor(&encoded, other_query, &key).is_err()); + + let mut tampered = URL_SAFE_NO_PAD.decode(&encoded).expect("cursor bytes"); + tampered[10] ^= 1; + assert!(decode_audit_cursor(&URL_SAFE_NO_PAD.encode(tampered), query, &key).is_err()); + } + + #[test] + fn admin_time_and_identity_helpers_are_bounded_and_domain_separated() { + assert_eq!(seconds_as_millis(1, false).expect("start"), 1_000); + assert_eq!(seconds_as_millis(1, true).expect("end"), 1_999); + assert!(seconds_as_millis(u64::MAX, false).is_err()); + let value = "caller-01"; + let operation = admin_identity(b"operation\0", value); + let correlation = admin_identity(b"correlation\0", value); + assert_ne!(operation, correlation); + assert!(!format!("{:?}", MycAuditCorrelationId::new(operation)).contains(value)); + } + + #[test] + fn admin_metric_keys_are_closed_and_bounded() { + for value in [ + "radroots_myc_service_ready", + "radroots_myc_service_phase_starting", + "radroots_myc_service_phase_degraded", + ] { + assert!(safe_metric_key(value)); + } + for value in ["", "Uppercase", "hyphen-key", "path/key", "secret:key"] { + assert!(!safe_metric_key(value)); + } + assert!(safe_metric_key(&"x".repeat(64))); + assert!(!safe_metric_key(&"x".repeat(65))); + } + + #[tokio::test] + async fn runtime_status_queries_real_connection_and_outbox_state() { + let directory = tempfile::tempdir().expect("temporary root"); + let runtime = crate::nip46_wave_080_a::runtime(directory.path()); + fs::create_dir_all(runtime.context().paths().state()).expect("state directory"); + fs::set_permissions( + runtime.context().paths().state(), + fs::Permissions::from_mode(0o700), + ) + .expect("state mode"); + let metadata = crate::nip46_wave_080_a::metadata(&runtime); + let configuration = crate::nip46_wave_080_a::configuration(); + let (applied_at, build) = crate::nip46_wave_080_a::migration_evidence(); + crate::initialize_myc_state(&runtime, &metadata, applied_at, &build) + .await + .expect("initialize"); + let host = crate::open_myc_state_read_write(&runtime, &metadata, applied_at, &build) + .await + .expect("host"); + let repository = host.repository(); + let empty = repository + .read_runtime_connection_counts() + .await + .expect("empty connection counts"); + assert_eq!(empty, MycConnectionCountsV1::default()); + let pending = crate::nip46_wave_080_a::admit_connect( + &repository, + &configuration, + 20, + "runtime-status-connect", + crate::nip46_wave_080_a::OBSERVED_AT_SECONDS, + crate::nip46_wave_080_a::RECEIVED_AT_MS, + 21, + crate::nip46_wave_080_a::RECEIVED_AT_MS + 1, + ) + .await; + assert!(pending.record().is_some()); + let counts = repository + .read_runtime_connection_counts() + .await + .expect("connection counts"); + assert_eq!(counts.pending(), 1); + assert_eq!(counts.active(), 0); + assert_eq!(counts.denied(), 0); + assert_eq!(counts.expired(), 0); + assert_eq!( + repository + .read_runtime_outbox_status() + .await + .expect("empty outbox"), + MycOutboxStatusV1::default() + ); + host.close().await.expect("close"); + + let inspection = crate::open_myc_state_inspection(&runtime, &metadata) + .await + .expect("inspection"); + let repository = inspection.repository(); + assert_eq!( + repository + .read_runtime_outbox_status() + .await + .expect("inspection outbox"), + MycOutboxStatusV1::default() + ); + assert_eq!( + repository + .read_runtime_connection_counts() + .await + .expect("inspection connection counts") + .pending(), + 1 + ); + inspection + .inspect_integrity( + radroots_service_sqlite::IntegrityCheckedAtUnixMs::new(1).expect("integrity time"), + ) + .await + .expect("inspection integrity"); + inspection.close().await.expect("inspection close"); + } +} diff --git a/src/runtime_nip46.rs b/src/runtime_nip46.rs @@ -0,0 +1,718 @@ +//! Transaction-free runtime coordination for one admitted NIP-46 event. + +use core::{fmt, str::FromStr as _}; +use std::{error::Error, sync::Arc}; + +use nostr::{JsonUtil as _, Kind, PublicKey as NostrPublicKey, Tag, Timestamp, UnsignedEvent}; +use radroots_identity::PublicKey; +use radroots_nostr_connect::{ + message::{RemoteSessionCapability, Response, SignedEvent as ConnectSignedEvent}, + permission::Permissions, + uri::RelayUrl, +}; +use radroots_service_host::{EntropySource, SystemEntropy, SystemWallClock, WallClock}; +use sha2::{Digest, Sha256}; + +use crate::{ + MycConfigDocumentV1, MycConnectionAdmissionPolicy, MycConnectionDecision, + MycConnectionDecisionRecord, MycConnectionNonce, MycConnectionPermissionSet, + MycConnectionPolicyGeneration, MycConnectionTimeUnixMs, MycNip46AdmissionLimits, + MycNip46AuthoredTimePolicy, MycNip46CommitRequest, MycNip46EncryptionContext, + MycNip46ObservedAtUnixSeconds, MycNip46ResponseCommitRequest, MycNip46Work, MycNip46WorkKind, + MycProviderCorrelationId, MycProviderDeadlineUnixMs, MycProviderNip44Version, + MycProviderOperation, MycProviderOperationId, MycProviderOperationInput, + MycProviderPublicIdentity, MycProviderResponseObservedAtUnixMs, MycProviderRole, + MycRateRelayId, MycRequestReceivedAtUnixMs, MycSignerOperationNonce, MycSignerRequestAdmission, + MycSignerRequestMethod, MycStateHost, MycTaskCancellation, admit_myc_nip46_event, + prepare_myc_nip46_decrypt_work, prepare_myc_nip46_request, prepare_myc_nip46_work, + provider_executor::MycProviderExecutor, verify_myc_nip46_event, +}; + +const RESPONSE_ENCRYPT_OPERATION_DOMAIN: &[u8] = + b"radroots.myc.nip46.response_encrypt.operation.v1\0"; +const RESPONSE_ENCRYPT_CORRELATION_DOMAIN: &[u8] = + b"radroots.myc.nip46.response_encrypt.correlation.v1\0"; +const RESPONSE_SIGN_OPERATION_DOMAIN: &[u8] = b"radroots.myc.nip46.response_sign.operation.v1\0"; +const RESPONSE_SIGN_CORRELATION_DOMAIN: &[u8] = + b"radroots.myc.nip46.response_sign.correlation.v1\0"; +const NIP46_RPC_KIND: u16 = 24_133; + +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub(crate) enum MycNip46DispatchDisposition { + Dropped, + PendingApproval, + Completed, + ExactCompletedReplay, +} + +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub(crate) enum MycNip46DispatchErrorKind { + Provider, + State, + Runtime, +} + +pub(crate) struct MycNip46DispatchError { + kind: MycNip46DispatchErrorKind, +} + +impl MycNip46DispatchError { + pub(crate) const fn kind(&self) -> MycNip46DispatchErrorKind { + self.kind + } +} + +impl fmt::Debug for MycNip46DispatchError { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter + .debug_struct("MycNip46DispatchError") + .field("kind", &self.kind) + .finish() + } +} + +impl fmt::Display for MycNip46DispatchError { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter.write_str("Myc NIP-46 dispatch failed") + } +} + +impl Error for MycNip46DispatchError {} + +const fn dispatch_error(kind: MycNip46DispatchErrorKind) -> MycNip46DispatchError { + MycNip46DispatchError { kind } +} + +pub(crate) struct MycRuntimeNip46Coordinator { + configuration: Arc<MycConfigDocumentV1>, + state: Arc<MycStateHost>, + providers: Arc<MycProviderExecutor>, + limits: MycNip46AdmissionLimits, + authored_time: MycNip46AuthoredTimePolicy, + provider_timeout_ms: u64, + response_relays: Box<[RelayUrl]>, +} + +#[derive(Clone, Copy, PartialEq, Eq)] +pub(crate) struct MycRuntimeNip46AdmissionEvidence { + request_nonce: [u8; 32], + received_at: MycRequestReceivedAtUnixMs, +} + +impl fmt::Debug for MycRuntimeNip46AdmissionEvidence { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter.write_str("MycRuntimeNip46AdmissionEvidence([redacted])") + } +} + +impl MycRuntimeNip46AdmissionEvidence { + pub(crate) fn new( + request_nonce: [u8; 32], + received_at_unix_ms: u64, + ) -> Result<Self, MycNip46DispatchError> { + Ok(Self { + request_nonce, + received_at: MycRequestReceivedAtUnixMs::new(received_at_unix_ms) + .map_err(|_| dispatch_error(MycNip46DispatchErrorKind::Runtime))?, + }) + } +} + +impl MycRuntimeNip46Coordinator { + pub(crate) fn new( + configuration: Arc<MycConfigDocumentV1>, + state: Arc<MycStateHost>, + providers: Arc<MycProviderExecutor>, + ) -> Result<Self, MycNip46DispatchError> { + let limits = MycNip46AdmissionLimits::from_config(&configuration) + .map_err(|_| dispatch_error(MycNip46DispatchErrorKind::Runtime))?; + let past = + configuration_integer(&configuration, "/transport/ingress/maximum_past_seconds")?; + let future = + configuration_integer(&configuration, "/transport/ingress/maximum_future_seconds")?; + let authored_time = MycNip46AuthoredTimePolicy::new(past, future) + .map_err(|_| dispatch_error(MycNip46DispatchErrorKind::Runtime))?; + let provider_timeout_ms = configuration_integer( + &configuration, + "/transport/ingress/subscription_deadline_ms", + )?; + let response_relays = response_relays(&configuration)?; + Ok(Self { + configuration, + state, + providers, + limits, + authored_time, + provider_timeout_ms, + response_relays, + }) + } + + pub(crate) async fn process( + &self, + raw_event: &[u8], + relay_id: MycRateRelayId, + observed_at_unix_ms: u64, + admission_evidence: MycRuntimeNip46AdmissionEvidence, + cancellation: &MycTaskCancellation, + ) -> Result<MycNip46DispatchDisposition, MycNip46DispatchError> { + let observed_seconds = observed_at_unix_ms / 1_000; + let bounded = match admit_myc_nip46_event(self.limits, raw_event) { + Ok(bounded) => bounded, + Err(_) => return Ok(MycNip46DispatchDisposition::Dropped), + }; + let observed = match MycNip46ObservedAtUnixSeconds::new(observed_seconds) { + Ok(observed) => observed, + Err(_) => return Ok(MycNip46DispatchDisposition::Dropped), + }; + let verified = match verify_myc_nip46_event( + bounded, + self.transport_binding()?, + observed, + self.authored_time, + ) { + Ok(verified) => verified, + Err(_) => return Ok(MycNip46DispatchDisposition::Dropped), + }; + let decrypt_deadline = self.provider_deadline()?; + let decrypt = match prepare_myc_nip46_decrypt_work( + verified, + self.transport_binding()?, + decrypt_deadline, + ) { + Ok(work) => work, + Err(_) => return Ok(MycNip46DispatchDisposition::Dropped), + }; + let decrypt_response = self + .providers + .execute( + decrypt.operation().owned_for_runtime(), + provider_observed_now()?, + cancellation, + ) + .await + .map_err(|_| dispatch_error(MycNip46DispatchErrorKind::Provider))?; + let decrypted = match decrypt.complete(decrypt_response, self.limits) { + Ok(request) => request, + Err(_) => return Ok(MycNip46DispatchDisposition::Dropped), + }; + let received_at = MycConnectionTimeUnixMs::new(admission_evidence.received_at.get()) + .map_err(|_| dispatch_error(MycNip46DispatchErrorKind::Runtime))?; + let prepared = match prepare_myc_nip46_request( + decrypted, + MycSignerOperationNonce::from_injected_entropy(admission_evidence.request_nonce), + admission_evidence.received_at, + ) { + Ok(request) => request, + Err(_) => return Ok(MycNip46DispatchDisposition::Dropped), + }; + let admission = self + .state + .repository() + .admit_signer_request(prepared.signer_request()) + .await + .map_err(|_| dispatch_error(MycNip46DispatchErrorKind::State))?; + if matches!(admission, MycSignerRequestAdmission::ConflictingReuse(_)) { + return Ok(MycNip46DispatchDisposition::Dropped); + } + if self + .state + .repository() + .read_nip46_response_by_operation(admission.record().operation_id()) + .await + .map_err(|_| dispatch_error(MycNip46DispatchErrorKind::State))? + .is_some() + { + return Ok(MycNip46DispatchDisposition::ExactCompletedReplay); + } + + let connection = if prepared.method() == MycSignerRequestMethod::Connect { + None + } else { + self.state + .repository() + .read_active_connection_for_client( + admission.record().client_public_key(), + received_at, + ) + .await + .map_err(|_| dispatch_error(MycNip46DispatchErrorKind::State))? + }; + let provider_deadline = matches!( + prepared.method(), + MycSignerRequestMethod::SignEvent + | MycSignerRequestMethod::Nip04Encrypt + | MycSignerRequestMethod::Nip04Decrypt + | MycSignerRequestMethod::Nip44Encrypt + | MycSignerRequestMethod::Nip44Decrypt + ) + .then(|| self.provider_deadline()) + .transpose()?; + let work = match prepare_myc_nip46_work( + prepared, + admission.record().clone(), + connection, + self.configuration.provider_contract(), + received_at, + provider_deadline, + ) { + Ok(work) => work, + Err(_) => return Ok(MycNip46DispatchDisposition::Dropped), + }; + + let connect_decision = if work.kind() == MycNip46WorkKind::Connect { + let request = self + .connection_request(&work, relay_id, received_at) + .await?; + let admission = self + .state + .repository() + .admit_connection(&request) + .await + .map_err(|_| dispatch_error(MycNip46DispatchErrorKind::State))?; + let Some(record) = admission.record().cloned() else { + return Ok(MycNip46DispatchDisposition::Dropped); + }; + if matches!( + record.decision(), + MycConnectionDecision::PendingApproval | MycConnectionDecision::Challenged + ) { + return Ok(MycNip46DispatchDisposition::PendingApproval); + } + Some(record) + } else { + None + }; + + let provider_response = if let Some(operation) = work.provider_operation() { + Some( + self.providers + .execute( + operation.owned_for_runtime(), + provider_observed_now()?, + cancellation, + ) + .await + .map_err(|_| dispatch_error(MycNip46DispatchErrorKind::Provider))?, + ) + } else { + None + }; + let completed_at = connection_time_now()?; + let completion = MycNip46CommitRequest::new( + &work, + connect_decision.as_ref(), + provider_response.as_ref(), + completed_at, + ) + .map_err(|_| dispatch_error(MycNip46DispatchErrorKind::Runtime))?; + let protocol_response = + self.protocol_response(&work, connect_decision.as_ref(), provider_response.as_ref())?; + let envelope = protocol_response + .into_envelope(work.request_record().request_id().as_str()) + .map_err(|_| dispatch_error(MycNip46DispatchErrorKind::Runtime))?; + let plaintext = serde_json::to_vec(&envelope) + .map_err(|_| dispatch_error(MycNip46DispatchErrorKind::Runtime))?; + let encrypted = self + .encrypt_response(&work, &plaintext, cancellation) + .await?; + let signed = self + .sign_response(&work, &encrypted, completed_at, cancellation) + .await?; + let commit = MycNip46ResponseCommitRequest::new( + &completion, + &signed.operation, + &signed.response, + crate::MycDeliveryTimeUnixMs::new(completed_at.get()) + .map_err(|_| dispatch_error(MycNip46DispatchErrorKind::Runtime))?, + ) + .map_err(|_| dispatch_error(MycNip46DispatchErrorKind::Runtime))?; + self.state + .repository() + .commit_nip46_response(&commit) + .await + .map_err(|_| dispatch_error(MycNip46DispatchErrorKind::State))?; + Ok(MycNip46DispatchDisposition::Completed) + } + + fn transport_binding(&self) -> Result<&crate::MycProviderBinding, MycNip46DispatchError> { + self.configuration + .provider_contract() + .binding(MycProviderRole::Transport) + .ok_or_else(|| dispatch_error(MycNip46DispatchErrorKind::Runtime)) + } + + fn provider_deadline(&self) -> Result<MycProviderDeadlineUnixMs, MycNip46DispatchError> { + wall_time_millis()? + .checked_add(self.provider_timeout_ms) + .and_then(|value| MycProviderDeadlineUnixMs::new(value).ok()) + .ok_or_else(|| dispatch_error(MycNip46DispatchErrorKind::Runtime)) + } + + async fn connection_request( + &self, + work: &MycNip46Work, + relay_id: MycRateRelayId, + observed_at: MycConnectionTimeUnixMs, + ) -> Result<crate::MycConnectionAdmissionRequest, MycNip46DispatchError> { + let client = work.request_record().client_public_key(); + let normalized = self.configuration.normalized(); + let contains = |pointer: &str| { + normalized + .pointer(pointer) + .and_then(serde_json::Value::as_array) + .is_some_and(|values| { + values + .iter() + .any(|value| value.as_str() == Some(client.as_hex())) + }) + }; + let policy = if contains("/policy/denied_clients") { + MycConnectionAdmissionPolicy::Denied + } else if contains("/policy/trusted_clients") { + MycConnectionAdmissionPolicy::Trusted + } else { + MycConnectionAdmissionPolicy::ExplicitApproval + }; + let authorized_until = if policy == MycConnectionAdmissionPolicy::Trusted + && normalized + .pointer("/policy/challenges/enabled") + .and_then(serde_json::Value::as_bool) + == Some(true) + { + let lifetime = configuration_integer( + &self.configuration, + "/policy/challenges/authorized_lifetime_ms", + )?; + Some( + observed_at + .get() + .checked_add(lifetime) + .and_then(|value| MycConnectionTimeUnixMs::new(value).ok()) + .ok_or_else(|| dispatch_error(MycNip46DispatchErrorKind::Runtime))?, + ) + } else { + None + }; + let generation = self + .state + .repository() + .current_configuration_generation() + .await + .map_err(|_| dispatch_error(MycNip46DispatchErrorKind::State))?; + let mut nonce = [0_u8; 32]; + SystemEntropy + .fill_bytes(&mut nonce) + .map_err(|_| dispatch_error(MycNip46DispatchErrorKind::Runtime))?; + work.connection_admission_request( + MycConnectionPolicyGeneration::new(u64::from(generation)) + .map_err(|_| dispatch_error(MycNip46DispatchErrorKind::Runtime))?, + MycConnectionNonce::from_injected_entropy(nonce), + observed_at, + authorized_until, + policy, + relay_id, + ) + .map_err(|_| dispatch_error(MycNip46DispatchErrorKind::Runtime)) + } + + fn protocol_response( + &self, + work: &MycNip46Work, + connect: Option<&MycConnectionDecisionRecord>, + provider: Option<&crate::MycVerifiedProviderResponse>, + ) -> Result<Response, MycNip46DispatchError> { + let user = PublicKey::from_hex( + self.configuration + .provider_contract() + .binding(MycProviderRole::User) + .ok_or_else(|| dispatch_error(MycNip46DispatchErrorKind::Runtime))? + .expected_identity() + .as_hex(), + ) + .map_err(|_| dispatch_error(MycNip46DispatchErrorKind::Runtime))?; + match work.method() { + MycSignerRequestMethod::Connect => { + match connect.map(MycConnectionDecisionRecord::decision) { + Some(MycConnectionDecision::Allowed) => Ok(Response::UserPublicKey(user)), + Some(MycConnectionDecision::Denied) => Ok(Response::Error { + result: None, + error: "connection_denied".to_owned(), + }), + _ => Err(dispatch_error(MycNip46DispatchErrorKind::Runtime)), + } + } + MycSignerRequestMethod::GetPublicKey => Ok(Response::UserPublicKey(user)), + MycSignerRequestMethod::GetSessionCapability => { + let permissions = protocol_permissions( + work.connection() + .ok_or_else(|| dispatch_error(MycNip46DispatchErrorKind::Runtime))? + .granted_permissions(), + )?; + Ok(Response::RemoteSessionCapability( + RemoteSessionCapability::try_new( + user, + self.response_relays.to_vec(), + permissions, + ) + .map_err(|_| dispatch_error(MycNip46DispatchErrorKind::Runtime))?, + )) + } + MycSignerRequestMethod::SignEvent => { + let bytes = provider + .and_then(crate::MycVerifiedProviderResponse::signed_event_bytes) + .ok_or_else(|| dispatch_error(MycNip46DispatchErrorKind::Runtime))?; + let json = core::str::from_utf8(bytes) + .map_err(|_| dispatch_error(MycNip46DispatchErrorKind::Runtime))?; + Ok(Response::SignedEvent( + ConnectSignedEvent::from_json(json) + .map_err(|_| dispatch_error(MycNip46DispatchErrorKind::Runtime))?, + )) + } + MycSignerRequestMethod::Nip04Encrypt => { + Ok(Response::Nip04Encrypt(protected_text(provider)?)) + } + MycSignerRequestMethod::Nip04Decrypt => { + Ok(Response::Nip04Decrypt(protected_text(provider)?)) + } + MycSignerRequestMethod::Nip44Encrypt => { + Ok(Response::Nip44Encrypt(protected_text(provider)?)) + } + MycSignerRequestMethod::Nip44Decrypt => { + Ok(Response::Nip44Decrypt(protected_text(provider)?)) + } + MycSignerRequestMethod::Ping => Ok(Response::Pong), + MycSignerRequestMethod::SwitchRelays => Ok(Response::RelayListUnchanged), + MycSignerRequestMethod::Logout => Ok(Response::LogoutAcknowledged), + } + } + + async fn encrypt_response( + &self, + work: &MycNip46Work, + plaintext: &[u8], + cancellation: &MycTaskCancellation, + ) -> Result<Vec<u8>, MycNip46DispatchError> { + let binding = self + .configuration + .provider_contract() + .binding(MycProviderRole::Transport) + .ok_or_else(|| dispatch_error(MycNip46DispatchErrorKind::Runtime))?; + let peer = + MycProviderPublicIdentity::new(work.request_record().client_public_key().as_hex()) + .map_err(|_| dispatch_error(MycNip46DispatchErrorKind::Runtime))?; + let input = match work.encryption_context() { + MycNip46EncryptionContext::Nip04 => { + MycProviderOperationInput::nip04_encrypt(peer, plaintext) + } + MycNip46EncryptionContext::Nip44V2 => MycProviderOperationInput::nip44_encrypt( + peer, + MycProviderNip44Version::V2, + plaintext, + ), + } + .map_err(|_| dispatch_error(MycNip46DispatchErrorKind::Runtime))?; + let operation = derived_operation( + binding, + work, + RESPONSE_ENCRYPT_OPERATION_DOMAIN, + RESPONSE_ENCRYPT_CORRELATION_DOMAIN, + self.provider_deadline()?, + input, + )?; + let response = self + .providers + .execute(operation, provider_observed_now()?, cancellation) + .await + .map_err(|_| dispatch_error(MycNip46DispatchErrorKind::Provider))?; + response + .protected_payload() + .map(<[u8]>::to_vec) + .ok_or_else(|| dispatch_error(MycNip46DispatchErrorKind::Provider)) + } + + async fn sign_response( + &self, + work: &MycNip46Work, + ciphertext: &[u8], + completed_at: MycConnectionTimeUnixMs, + cancellation: &MycTaskCancellation, + ) -> Result<SignedRuntimeResponse, MycNip46DispatchError> { + let content = core::str::from_utf8(ciphertext) + .map_err(|_| dispatch_error(MycNip46DispatchErrorKind::Provider))?; + let user_binding = self + .configuration + .provider_contract() + .binding(MycProviderRole::User) + .ok_or_else(|| dispatch_error(MycNip46DispatchErrorKind::Runtime))?; + let user = NostrPublicKey::from_hex(user_binding.expected_identity().as_hex()) + .map_err(|_| dispatch_error(MycNip46DispatchErrorKind::Runtime))?; + let client = NostrPublicKey::from_hex(work.request_record().client_public_key().as_hex()) + .map_err(|_| dispatch_error(MycNip46DispatchErrorKind::Runtime))?; + let unsigned = UnsignedEvent::new( + user, + Timestamp::from_secs(completed_at.get() / 1_000), + Kind::Custom(NIP46_RPC_KIND), + vec![Tag::public_key(client)], + content, + ); + let input = MycProviderOperationInput::sign_event(unsigned.as_json().as_bytes()) + .map_err(|_| dispatch_error(MycNip46DispatchErrorKind::Runtime))?; + let operation = derived_operation( + user_binding, + work, + RESPONSE_SIGN_OPERATION_DOMAIN, + RESPONSE_SIGN_CORRELATION_DOMAIN, + self.provider_deadline()?, + input, + )?; + let response = self + .providers + .execute( + operation.owned_for_runtime(), + provider_observed_now()?, + cancellation, + ) + .await + .map_err(|_| dispatch_error(MycNip46DispatchErrorKind::Provider))?; + Ok(SignedRuntimeResponse { + operation, + response, + }) + } +} + +struct SignedRuntimeResponse { + operation: MycProviderOperation, + response: crate::MycVerifiedProviderResponse, +} + +fn derived_operation( + binding: &crate::MycProviderBinding, + work: &MycNip46Work, + operation_domain: &[u8], + correlation_domain: &[u8], + deadline: MycProviderDeadlineUnixMs, + input: MycProviderOperationInput, +) -> Result<MycProviderOperation, MycNip46DispatchError> { + let identity = work.request_record().operation_id(); + MycProviderOperation::new( + binding, + MycProviderOperationId::from_bytes(derived_identifier( + operation_domain, + identity.as_bytes(), + )), + MycProviderCorrelationId::from_bytes(derived_identifier( + correlation_domain, + identity.as_bytes(), + )), + deadline, + input, + ) + .map_err(|_| dispatch_error(MycNip46DispatchErrorKind::Runtime)) +} + +fn derived_identifier(domain: &[u8], operation_id: &[u8; 32]) -> [u8; 32] { + let mut hasher = Sha256::new(); + hasher.update(domain); + hasher.update(operation_id); + hasher.finalize().into() +} + +fn protected_text( + provider: Option<&crate::MycVerifiedProviderResponse>, +) -> Result<String, MycNip46DispatchError> { + provider + .and_then(crate::MycVerifiedProviderResponse::protected_payload) + .and_then(|bytes| core::str::from_utf8(bytes).ok()) + .map(str::to_owned) + .ok_or_else(|| dispatch_error(MycNip46DispatchErrorKind::Runtime)) +} + +fn protocol_permissions( + permissions: &MycConnectionPermissionSet, +) -> Result<Permissions, MycNip46DispatchError> { + let value = permissions + .permissions() + .iter() + .map(|permission| permission.code()) + .collect::<Vec<_>>() + .join(","); + Permissions::from_str(&value).map_err(|_| dispatch_error(MycNip46DispatchErrorKind::Runtime)) +} + +fn response_relays( + configuration: &MycConfigDocumentV1, +) -> Result<Box<[RelayUrl]>, MycNip46DispatchError> { + let relays = configuration + .normalized() + .pointer("/relays") + .and_then(serde_json::Value::as_array) + .ok_or_else(|| dispatch_error(MycNip46DispatchErrorKind::Runtime))?; + relays + .iter() + .filter(|relay| relay.pointer("/write").and_then(serde_json::Value::as_bool) == Some(true)) + .map(|relay| { + relay + .pointer("/url") + .and_then(serde_json::Value::as_str) + .ok_or_else(|| dispatch_error(MycNip46DispatchErrorKind::Runtime)) + .and_then(|url| { + RelayUrl::parse(url) + .map_err(|_| dispatch_error(MycNip46DispatchErrorKind::Runtime)) + }) + }) + .collect::<Result<Vec<_>, _>>() + .map(Vec::into_boxed_slice) +} + +fn wall_time_millis() -> Result<u64, MycNip46DispatchError> { + SystemWallClock + .now_utc() + .map_err(|_| dispatch_error(MycNip46DispatchErrorKind::Runtime))? + .get() + .checked_mul(1_000) + .filter(|value| i64::try_from(*value).is_ok()) + .ok_or_else(|| dispatch_error(MycNip46DispatchErrorKind::Runtime)) +} + +fn connection_time_now() -> Result<MycConnectionTimeUnixMs, MycNip46DispatchError> { + MycConnectionTimeUnixMs::new(wall_time_millis()?) + .map_err(|_| dispatch_error(MycNip46DispatchErrorKind::Runtime)) +} + +fn provider_observed_now() -> Result<MycProviderResponseObservedAtUnixMs, MycNip46DispatchError> { + MycProviderResponseObservedAtUnixMs::new(wall_time_millis()?) + .map_err(|_| dispatch_error(MycNip46DispatchErrorKind::Runtime)) +} + +fn configuration_integer( + configuration: &MycConfigDocumentV1, + pointer: &str, +) -> Result<u64, MycNip46DispatchError> { + configuration + .normalized() + .pointer(pointer) + .and_then(serde_json::Value::as_u64) + .ok_or_else(|| dispatch_error(MycNip46DispatchErrorKind::Runtime)) +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn admission_evidence_is_copyable_across_exact_runtime_retries() { + let evidence = MycRuntimeNip46AdmissionEvidence::new([0x5a; 32], 1_725_000_000_000) + .expect("admission evidence"); + assert_eq!(evidence, evidence); + assert_eq!(evidence.request_nonce, [0x5a; 32]); + assert_eq!(evidence.received_at.get(), 1_725_000_000_000); + assert_eq!( + format!("{evidence:?}"), + "MycRuntimeNip46AdmissionEvidence([redacted])" + ); + assert!(MycRuntimeNip46AdmissionEvidence::new([0x5a; 32], 0).is_err()); + assert!(MycRuntimeNip46AdmissionEvidence::new([0x5a; 32], u64::MAX).is_err()); + } +} diff --git a/src/runtime_signal.rs b/src/runtime_signal.rs @@ -0,0 +1,50 @@ +//! Binary-owned process-signal injection for the Myc daemon. + +use core::{fmt, future::Future, pin::Pin}; + +/// One normalized process signal supplied by the Myc binary. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub enum MycProcessSignal { + Interrupt, + #[cfg(unix)] + Terminate, +} + +impl MycProcessSignal { + #[must_use] + pub const fn as_str(self) -> &'static str { + match self { + Self::Interrupt => "interrupt", + #[cfg(unix)] + Self::Terminate => "terminate", + } + } +} + +impl fmt::Display for MycProcessSignal { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter.write_str(self.as_str()) + } +} + +/// Boxed wait returned by a binary-owned signal source. +pub type MycProcessSignalFuture<'a> = + Pin<Box<dyn Future<Output = Option<MycProcessSignal>> + Send + 'a>>; + +/// Process-signal source installed only by the executable boundary. +pub trait MycProcessSignalSource: Send { + fn next_signal(&mut self) -> MycProcessSignalFuture<'_>; +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn normalized_signal_names_are_stable() { + assert_eq!(MycProcessSignal::Interrupt.as_str(), "interrupt"); + assert_eq!(MycProcessSignal::Interrupt.to_string(), "interrupt"); + #[cfg(unix)] + assert_eq!(MycProcessSignal::Terminate.as_str(), "terminate"); + } +} diff --git a/src/runtime_supervision.rs b/src/runtime_supervision.rs @@ -32,6 +32,10 @@ pub struct MycTaskCancellation { } impl MycTaskCancellation { + pub(crate) const fn from_host(inner: CancellationToken) -> Self { + Self { inner } + } + #[cfg(any(target_os = "linux", target_os = "macos"))] pub(crate) fn uncancelled() -> Self { Self { diff --git a/src/state_admin.rs b/src/state_admin.rs @@ -1,6 +1,6 @@ //! Bounded durable idempotency for permissioned Myc admin mutations. -use core::fmt; +use core::{fmt, future::Future, pin::Pin}; use std::error::Error; use radroots_service_sqlite::{ @@ -292,6 +292,133 @@ pub enum MycAdminOperationCompletion { } impl MycStateRepository<'_> { + /// Atomically commits one SQLite-only admin mutation and its replay receipt. + /// + /// The supplied operation runs inside the same governed transaction that + /// admits the operation identifier and records the canonical response. A + /// domain failure therefore leaves neither a prepared journal row nor a + /// partial domain effect. + pub(crate) async fn execute_database_admin_operation<F>( + &self, + request: &MycAdminRequestDocument, + completed_at: MycAdminOperationTimeUnixMs, + policy: MycAdminOperationJournalPolicy, + operation: F, + ) -> Result<MycAdminResponseDocument, MycAdminOperationError> + where + F: for<'a, 'b> FnOnce( + &'a mut ServiceSqliteTransaction<'b>, + ) -> AdminDatabaseOperationFuture<'a> + + Send + + 'static, + { + if !self.is_writable() { + return Err(MycAdminOperationError::new( + MycAdminOperationErrorKind::InvalidMode, + )); + } + let binding = AdminRequestBinding::from_document(request)?; + let expires_at = completed_at + .get() + .checked_add(policy.completed_retention_ms()) + .filter(|value| i64::try_from(*value).is_ok()) + .ok_or_else(|| MycAdminOperationError::new(MycAdminOperationErrorKind::InvalidInput))?; + let expected = PersistedMetadata::from(self.expected()); + self.host() + .transaction(move |transaction| { + Box::pin(async move { + require_expected_metadata(transaction, &expected) + .await + .map_err(AdminJournalOperationError::from)?; + let prepared = match prepare_operation(transaction, &binding, completed_at) + .await? + { + MycAdminOperationAdmission::ExactReplay(response) => return Ok(response), + MycAdminOperationAdmission::Prepared(prepared) => prepared, + }; + let response = operation(transaction).await?; + if response.route() != binding.route + || response.canonical_bytes().is_empty() + || response.canonical_bytes().len() + > MYC_ADMIN_OPERATION_RESPONSE_MODEL_MAX_BYTES + { + return Err(AdminJournalOperationError::InvalidInput); + } + complete_operation( + transaction, + &PreparedBinding::from_prepared(&prepared), + response.canonical_bytes(), + completed_at, + expires_at, + ) + .await?; + Ok(response) + }) + }) + .await + .map_err(map_transaction_error) + } + + pub(crate) async fn complete_prepared_database_admin_operation<F>( + &self, + prepared: &MycPreparedAdminOperation, + completed_at: MycAdminOperationTimeUnixMs, + policy: MycAdminOperationJournalPolicy, + operation: F, + ) -> Result<MycAdminResponseDocument, MycAdminOperationError> + where + F: for<'a, 'b> FnOnce( + &'a mut ServiceSqliteTransaction<'b>, + ) -> AdminDatabaseOperationFuture<'a> + + Send + + 'static, + { + if !self.is_writable() { + return Err(MycAdminOperationError::new( + MycAdminOperationErrorKind::InvalidMode, + )); + } + if completed_at < prepared.prepared_at { + return Err(MycAdminOperationError::new( + MycAdminOperationErrorKind::InvalidInput, + )); + } + let expires_at = completed_at + .get() + .checked_add(policy.completed_retention_ms()) + .filter(|value| i64::try_from(*value).is_ok()) + .ok_or_else(|| MycAdminOperationError::new(MycAdminOperationErrorKind::InvalidInput))?; + let binding = PreparedBinding::from_prepared(prepared); + let expected = PersistedMetadata::from(self.expected()); + self.host() + .transaction(move |transaction| { + Box::pin(async move { + require_expected_metadata(transaction, &expected) + .await + .map_err(AdminJournalOperationError::from)?; + let response = operation(transaction).await?; + if response.route() != binding.route + || response.canonical_bytes().is_empty() + || response.canonical_bytes().len() + > MYC_ADMIN_OPERATION_RESPONSE_MODEL_MAX_BYTES + { + return Err(AdminJournalOperationError::InvalidInput); + } + complete_operation( + transaction, + &binding, + response.canonical_bytes(), + completed_at, + expires_at, + ) + .await?; + Ok(response) + }) + }) + .await + .map_err(map_transaction_error) + } + /// Prunes a bounded expired prefix and admits or replays one mutation. pub async fn prepare_admin_operation( &self, @@ -423,7 +550,7 @@ enum StoredOperation { } #[derive(Clone, Copy, Debug, PartialEq, Eq)] -enum AdminJournalOperationError { +pub(crate) enum AdminJournalOperationError { InvalidInput, Conflict, OutcomeUnknown, @@ -432,6 +559,14 @@ enum AdminJournalOperationError { Storage, } +pub(crate) type AdminDatabaseOperationFuture<'a> = Pin< + Box< + dyn Future<Output = Result<MycAdminResponseDocument, AdminJournalOperationError>> + + Send + + 'a, + >, +>; + impl From<RepositoryOperationError> for AdminJournalOperationError { fn from(error: RepositoryOperationError) -> Self { match error { @@ -1065,6 +1200,117 @@ mod tests { } #[tokio::test] + async fn database_only_admin_effect_and_receipt_share_one_transaction() { + let (_directory, _runtime, _metadata, host) = fixture().await; + let repository = host.repository(); + let request = request("atomic-admin-1", "connection-1", 1); + let error = repository + .execute_database_admin_operation( + &request, + MycAdminOperationTimeUnixMs::new(100).expect("time"), + MycAdminOperationJournalPolicy::seven_days(), + |transaction| { + Box::pin(async move { + sqlx::query( + r#"INSERT INTO connection_rate_windows ( + rate_kind, subject_scope, subject_sha256, + window_started_at_unix_ms, window_ends_at_unix_ms, + accepted_count, rejected_count, lifetime_accepted_count, + lifetime_rejected_count, last_observed_at_unix_ms, + retention_expires_at_unix_ms + ) VALUES ('connection_admission', 'global', ?, 1, 2, 0, 0, 0, 0, 1, 2)"#, + ) + .bind([0xaa_u8; 32].as_slice()) + .execute(&mut *transaction) + .await + .map_err(|_| AdminJournalOperationError::Storage)?; + Err(AdminJournalOperationError::Binding) + }) + }, + ) + .await + .expect_err("domain failure rolls back"); + assert_eq!(error.kind(), MycAdminOperationErrorKind::Binding); + assert_eq!( + repository + .host() + .transaction(|transaction| { + Box::pin(async move { + sqlx::query_scalar::<_, i64>("SELECT COUNT(*) FROM connection_rate_windows") + .fetch_one(&mut *transaction) + .await + }) + }) + .await + .expect("rolled-back effect"), + 0 + ); + + let expected = response_document("atomic-admin-1", 1); + let expected_bytes = expected.canonical_bytes().to_vec(); + let committed = repository + .execute_database_admin_operation( + &request, + MycAdminOperationTimeUnixMs::new(101).expect("time"), + MycAdminOperationJournalPolicy::seven_days(), + move |transaction| { + let expected_bytes = expected_bytes.clone(); + Box::pin(async move { + sqlx::query( + r#"INSERT INTO connection_rate_windows ( + rate_kind, subject_scope, subject_sha256, + window_started_at_unix_ms, window_ends_at_unix_ms, + accepted_count, rejected_count, lifetime_accepted_count, + lifetime_rejected_count, last_observed_at_unix_ms, + retention_expires_at_unix_ms + ) VALUES ('connection_admission', 'global', ?, 3, 4, 0, 0, 0, 0, 3, 4)"#, + ) + .bind([0xbb_u8; 32].as_slice()) + .execute(&mut *transaction) + .await + .map_err(|_| AdminJournalOperationError::Storage)?; + MycAdminResponseDocument::from_canonical_bytes( + MycAdminRoute::ConnectionApprove, + &expected_bytes, + ) + .map_err(|_| AdminJournalOperationError::Binding) + }) + }, + ) + .await + .expect("atomic commit"); + let replay = repository + .execute_database_admin_operation( + &request, + MycAdminOperationTimeUnixMs::new(102).expect("time"), + MycAdminOperationJournalPolicy::seven_days(), + |_transaction| { + Box::pin( + async move { panic!("exact replay must not execute the domain operation") }, + ) + }, + ) + .await + .expect("exact replay"); + assert_eq!(replay.canonical_bytes(), committed.canonical_bytes()); + assert_eq!( + repository + .host() + .transaction(|transaction| { + Box::pin(async move { + sqlx::query_scalar::<_, i64>("SELECT COUNT(*) FROM connection_rate_windows") + .fetch_one(&mut *transaction) + .await + }) + }) + .await + .expect("single committed effect"), + 1 + ); + host.close().await.expect("close"); + } + + #[tokio::test] async fn prepared_backup_shape_is_ambiguous_and_persists_no_path_or_request_content() { let (_directory, _runtime, _metadata, host) = fixture().await; let request = MycAdminRequestDocument::mutation_for_test( diff --git a/src/state_connection.rs b/src/state_connection.rs @@ -10,6 +10,8 @@ use sha2::{Digest, Sha256}; use sqlx::Row; use url::{Host, Url}; +#[cfg(any(target_os = "linux", target_os = "macos"))] +use crate::state_admin::AdminJournalOperationError; use crate::state_governance::{ AuditEvidence, GovernanceOperationError, MycAuditCorrelationId, MycAuditKind, MycAuditOutcome, MycAuditReasonCode, MycRateLimitClass, MycRateLimitPolicy, MycRateRelayId, connection_subject, @@ -101,6 +103,35 @@ FROM connections WHERE connection_id = ? LIMIT 2"#; +const READ_ACTIVE_CONNECTION_FOR_CLIENT_SQL: &str = r#"SELECT + CASE WHEN typeof(connection_id) = 'blob' AND length(connection_id) = 32 + THEN connection_id ELSE NULL END AS connection_id +FROM connections +WHERE client_public_key = ? AND status = 'active' + AND (authorized_until_unix_ms IS NULL OR authorized_until_unix_ms >= ?) +ORDER BY updated_at_unix_ms DESC, connection_id ASC +LIMIT 2"#; + +const READ_RUNTIME_CONNECTION_COUNTS_SQL: &str = r#"SELECT + CASE WHEN typeof(status) = 'text' AND length(CAST(status AS BLOB)) <= 16 + THEN status ELSE NULL END AS status, + COUNT(*) AS row_count +FROM connections +GROUP BY status +LIMIT 5"#; + +const READ_ADMIN_CONNECTION_PAGE_SQL: &str = r#"SELECT + CASE WHEN typeof(connection_id) = 'blob' AND length(connection_id) = 32 + THEN connection_id ELSE NULL END AS connection_id, + updated_at_unix_ms +FROM connections +WHERE updated_at_unix_ms <= ? + AND (? IS NULL OR status = ?) + AND (? IS NULL OR updated_at_unix_ms < ? + OR (updated_at_unix_ms = ? AND connection_id > ?)) +ORDER BY updated_at_unix_ms DESC, connection_id ASC +LIMIT ?"#; + const READ_PERMISSIONS_SQL: &str = r#"SELECT CASE WHEN typeof(permission_code) = 'text' AND length(CAST(permission_code AS BLOB)) BETWEEN 1 AND 64 @@ -122,6 +153,18 @@ SET status = 'expired', updated_at_unix_ms = ?, authorized_until_unix_ms = NULL WHERE connection_id = ? AND status = 'active' AND policy_generation = ? AND authorized_until_unix_ms IS NOT NULL AND authorized_until_unix_ms < ?"#; +const REVOKE_CONNECTION_SQL: &str = r#"UPDATE connections +SET status = 'expired', updated_at_unix_ms = ?, authorized_until_unix_ms = NULL +WHERE connection_id = ? AND status = 'active' AND policy_generation = ?"#; + +const READ_PENDING_DECISION_FOR_CONNECTION_SQL: &str = r#"SELECT + CASE WHEN typeof(operation_id) = 'blob' AND length(operation_id) = 32 + THEN operation_id ELSE NULL END AS operation_id +FROM nip46_request_decisions +WHERE connection_id = ? AND decision = 'pending_approval' AND policy_generation = ? +ORDER BY decided_at_unix_ms DESC, operation_id ASC +LIMIT 2"#; + const UPDATE_APPROVAL_DECISION_SQL: &str = r#"UPDATE nip46_request_decisions SET decision = ?, reason_code = ?, decided_at_unix_ms = ? WHERE operation_id = ? AND connection_id = ? AND decision = 'pending_approval' @@ -153,6 +196,13 @@ FROM connection_auth_challenges WHERE operation_id = ? LIMIT 2"#; +const READ_CHALLENGE_OPERATION_BY_ID_SQL: &str = r#"SELECT + CASE WHEN typeof(operation_id) = 'blob' AND length(operation_id) = 32 + THEN operation_id ELSE NULL END AS operation_id +FROM connection_auth_challenges +WHERE challenge_id = ? +LIMIT 2"#; + const RESOLVE_CHALLENGE_SQL: &str = r#"UPDATE connection_auth_challenges SET state = ?, resolved_at_unix_ms = ? WHERE challenge_id = ? AND connection_id = ? AND operation_id = ? @@ -258,7 +308,7 @@ pub enum MycConnectionPermission { } impl MycConnectionPermission { - fn code(self) -> String { + pub(crate) fn code(self) -> String { match self { Self::GetPublicKey => "get_public_key".into(), Self::GetSessionCapability => "get_session_capability".into(), @@ -392,6 +442,10 @@ macro_rules! digest_id { pub struct $name([u8; 32]); impl $name { + pub(crate) const fn from_bytes(bytes: [u8; 32]) -> Self { + Self(bytes) + } + /// Returns the exact stable identity bytes. #[must_use] pub const fn as_bytes(&self) -> &[u8; 32] { @@ -640,6 +694,15 @@ impl MycConnectionStatus { _ => None, } } + + pub(crate) const fn admin_state(self) -> &'static str { + match self { + Self::Pending => "pending", + Self::Active => "approved", + Self::Denied => "rejected", + Self::Expired => "revoked", + } + } } /// Validated durable connection record. @@ -656,6 +719,21 @@ pub struct MycConnectionRecord { authorized_until: Option<MycConnectionTimeUnixMs>, } +pub(crate) struct MycAdminConnectionPage { + items: Box<[MycConnectionRecord]>, + next: Option<(MycConnectionTimeUnixMs, MycConnectionId)>, +} + +impl MycAdminConnectionPage { + pub(crate) fn items(&self) -> &[MycConnectionRecord] { + &self.items + } + + pub(crate) const fn next(&self) -> Option<(MycConnectionTimeUnixMs, MycConnectionId)> { + self.next + } +} + impl MycConnectionRecord { #[must_use] pub const fn id(&self) -> MycConnectionId { @@ -694,6 +772,15 @@ impl MycConnectionRecord { self.authorized_until } + pub(crate) fn admin_permissions(&self) -> String { + self.granted_permissions + .permissions() + .iter() + .map(|permission| permission.code()) + .collect::<Vec<_>>() + .join(",") + } + #[cfg(test)] pub(crate) fn active_for_test( client_public_key: MycNip46ClientPublicKey, @@ -806,6 +893,16 @@ pub enum MycConnectionOperatorDecision { Deny, } +#[cfg(any(target_os = "linux", target_os = "macos"))] +pub(crate) enum MycAdminConnectionAction { + Approve { + permissions: MycConnectionPermissionSet, + authorized_until: Option<MycConnectionTimeUnixMs>, + }, + Reject, + Revoke, +} + impl fmt::Debug for MycConnectionOperatorDecision { fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { formatter.write_str(match self { @@ -1079,6 +1176,112 @@ impl fmt::Debug for MycAuthorizationChallengeAuthorization { } impl MycStateRepository<'_> { + pub(crate) async fn read_runtime_connection_counts( + &self, + ) -> Result<crate::MycConnectionCountsV1, MycStateRepositoryError> { + let expected = PersistedMetadata::from(self.expected()); + self.host() + .transaction(move |transaction| { + Box::pin(async move { + verify_metadata(transaction, &expected).await?; + let rows = sqlx::query(READ_RUNTIME_CONNECTION_COUNTS_SQL) + .fetch_all(&mut *transaction) + .await + .map_err(|_| ConnectionOperationError::Storage)?; + if rows.len() > 4 { + return Err(ConnectionOperationError::Binding); + } + let mut counts = [None; 4]; + for row in rows { + let status = bounded_text(&row, "status")?; + let index = match MycConnectionStatus::parse(status) { + Some(MycConnectionStatus::Pending) => 0, + Some(MycConnectionStatus::Active) => 1, + Some(MycConnectionStatus::Denied) => 2, + Some(MycConnectionStatus::Expired) => 3, + None => return Err(ConnectionOperationError::Binding), + }; + let value = row + .try_get::<i64, _>("row_count") + .ok() + .and_then(|value| u64::try_from(value).ok()) + .ok_or(ConnectionOperationError::Binding)?; + if counts[index].replace(value).is_some() { + return Err(ConnectionOperationError::Binding); + } + } + Ok(crate::MycConnectionCountsV1::new( + counts[0].unwrap_or(0), + counts[1].unwrap_or(0), + counts[2].unwrap_or(0), + counts[3].unwrap_or(0), + )) + }) + }) + .await + .map_err(map_transaction_error) + } + + pub(crate) async fn read_admin_connection_page( + &self, + limit: u16, + status: Option<MycConnectionStatus>, + snapshot: MycConnectionTimeUnixMs, + before: Option<(MycConnectionTimeUnixMs, MycConnectionId)>, + ) -> Result<MycAdminConnectionPage, MycStateRepositoryError> { + if limit == 0 || limit > 200 { + return Err(MycStateRepositoryError::new( + MycStateRepositoryErrorKind::Binding, + )); + } + let expected = PersistedMetadata::from(self.expected()); + self.host() + .transaction(move |transaction| { + Box::pin(async move { + verify_metadata(transaction, &expected).await?; + read_admin_connection_page(transaction, limit, status, snapshot, before).await + }) + }) + .await + .map_err(map_transaction_error) + } + + /// Resolves the sole active session for one verified transport client. + /// + /// Multiple simultaneously active sessions are ambiguous because NIP-46 + /// requests carry no server connection identifier; fail closed rather + /// than silently selecting a different authorization grant. + pub(crate) async fn read_active_connection_for_client( + &self, + client: &MycNip46ClientPublicKey, + observed_at: MycConnectionTimeUnixMs, + ) -> Result<Option<MycConnectionRecord>, MycStateRepositoryError> { + let expected = PersistedMetadata::from(self.expected()); + let client = client.clone(); + self.host() + .transaction(move |transaction| { + Box::pin(async move { + verify_metadata(transaction, &expected).await?; + let rows = sqlx::query(READ_ACTIVE_CONNECTION_FOR_CLIENT_SQL) + .bind(client.as_hex()) + .bind(observed_at.sqlite_value()) + .fetch_all(&mut *transaction) + .await + .map_err(|_| ConnectionOperationError::Storage)?; + match rows.as_slice() { + [] => Ok(None), + [row] => { + let id = MycConnectionId(exact_digest(row, "connection_id")?); + read_connection(transaction, id).await.map(Some) + } + _ => Err(ConnectionOperationError::Binding), + } + }) + }) + .await + .map_err(map_transaction_error) + } + /// Reads one fully validated durable connection decision by stable operation identity. pub async fn read_connection_decision( &self, @@ -1294,6 +1497,47 @@ impl MycStateRepository<'_> { } } +async fn read_admin_connection_page( + transaction: &mut ServiceSqliteTransaction<'_>, + limit: u16, + status: Option<MycConnectionStatus>, + snapshot: MycConnectionTimeUnixMs, + before: Option<(MycConnectionTimeUnixMs, MycConnectionId)>, +) -> Result<MycAdminConnectionPage, ConnectionOperationError> { + let status = status.map(MycConnectionStatus::as_str); + let before_time = before.map(|(time, _)| time.sqlite_value()); + let before_id = before.map(|(_, id)| id.as_bytes().to_vec()); + let fetch_limit = i64::from(limit) + 1; + let rows = sqlx::query(READ_ADMIN_CONNECTION_PAGE_SQL) + .bind(snapshot.sqlite_value()) + .bind(status) + .bind(status) + .bind(before_time) + .bind(before_time) + .bind(before_time) + .bind(before_id) + .bind(fetch_limit) + .fetch_all(&mut *transaction) + .await + .map_err(|_| ConnectionOperationError::Storage)?; + if rows.len() > usize::try_from(fetch_limit).map_err(|_| ConnectionOperationError::Binding)? { + return Err(ConnectionOperationError::Binding); + } + let has_more = rows.len() > usize::from(limit); + let mut items = Vec::with_capacity(rows.len().min(usize::from(limit))); + for row in rows.iter().take(usize::from(limit)) { + let id = MycConnectionId(exact_digest(row, "connection_id")?); + items.push(read_connection(transaction, id).await?); + } + let next = has_more + .then(|| items.last().map(|item| (item.updated_at, item.id))) + .flatten(); + Ok(MycAdminConnectionPage { + items: items.into_boxed_slice(), + next, + }) +} + async fn record_connection_expiry_audit( transaction: &mut ServiceSqliteTransaction<'_>, correlation: MycAuditCorrelationId, @@ -1321,6 +1565,174 @@ enum ConnectionOperationError { Storage, } +#[cfg(any(target_os = "linux", target_os = "macos"))] +impl From<ConnectionOperationError> for AdminJournalOperationError { + fn from(error: ConnectionOperationError) -> Self { + match error { + ConnectionOperationError::Binding => Self::Binding, + ConnectionOperationError::Storage => Self::Storage, + } + } +} + +#[cfg(any(target_os = "linux", target_os = "macos"))] +pub(crate) async fn apply_admin_connection_mutation( + transaction: &mut ServiceSqliteTransaction<'_>, + connection_id: MycConnectionId, + policy_generation: MycConnectionPolicyGeneration, + observed_at: MycConnectionTimeUnixMs, + audit_correlation: MycAuditCorrelationId, + action: MycAdminConnectionAction, +) -> Result<(MycConnectionRecord, MycConnectionRecord), AdminJournalOperationError> { + let before = read_connection(transaction, connection_id).await?; + if before.policy_generation != policy_generation { + return Err(AdminJournalOperationError::Conflict); + } + let after = match action { + MycAdminConnectionAction::Approve { + permissions, + authorized_until, + } => { + let operation_id = + read_pending_decision_operation(transaction, connection_id, policy_generation) + .await?; + decide_connection( + transaction, + operation_id, + connection_id, + policy_generation, + observed_at, + audit_correlation, + MycConnectionOperatorDecision::Approve { + granted_permissions: permissions, + authorized_until, + }, + ) + .await? + } + MycAdminConnectionAction::Reject => { + let operation_id = + read_pending_decision_operation(transaction, connection_id, policy_generation) + .await?; + decide_connection( + transaction, + operation_id, + connection_id, + policy_generation, + observed_at, + audit_correlation, + MycConnectionOperatorDecision::Deny, + ) + .await? + } + MycAdminConnectionAction::Revoke => { + if before.status == MycConnectionStatus::Expired { + before.clone() + } else { + if before.status != MycConnectionStatus::Active || observed_at < before.updated_at { + return Err(AdminJournalOperationError::Conflict); + } + let result = sqlx::query(REVOKE_CONNECTION_SQL) + .bind(observed_at.sqlite_value()) + .bind(connection_id.as_bytes().as_slice()) + .bind(policy_generation.sqlite_value()) + .execute(&mut *transaction) + .await + .map_err(|_| AdminJournalOperationError::Storage)?; + require_one(result.rows_affected())?; + let record = read_connection(transaction, connection_id).await?; + record_operator_audit( + transaction, + audit_correlation, + observed_at, + MycAuditOutcome::Succeeded, + MycAuditReasonCode::ConnectionExpired, + ) + .await?; + record + } + } + }; + Ok((before, after)) +} + +#[cfg(any(target_os = "linux", target_os = "macos"))] +pub(crate) async fn apply_admin_challenge_require( + transaction: &mut ServiceSqliteTransaction<'_>, + request: MycAuthorizationChallengeRequest, + rate_policy: MycRateLimitPolicy, +) -> Result<MycAuthorizationChallengeRecord, AdminJournalOperationError> { + match issue_challenge(transaction, &request, rate_policy).await? { + MycAuthorizationChallengeAdmission::Created(record) + | MycAuthorizationChallengeAdmission::ExactReplay(record) => Ok(record), + MycAuthorizationChallengeAdmission::RateLimited => { + Err(AdminJournalOperationError::ResourceExhausted) + } + } +} + +#[cfg(any(target_os = "linux", target_os = "macos"))] +pub(crate) async fn apply_admin_challenge_authorize( + transaction: &mut ServiceSqliteTransaction<'_>, + challenge_id: MycAuthorizationChallengeId, + policy_generation: MycConnectionPolicyGeneration, + observed_at: MycConnectionTimeUnixMs, + rate_policy: MycRateLimitPolicy, +) -> Result<MycAuthorizationChallengeRecord, AdminJournalOperationError> { + let rows = sqlx::query(READ_CHALLENGE_OPERATION_BY_ID_SQL) + .bind(challenge_id.as_bytes().as_slice()) + .fetch_all(&mut *transaction) + .await + .map_err(|_| AdminJournalOperationError::Storage)?; + let operation_id = match rows.as_slice() { + [row] => MycSignerOperationId::from_persisted(exact_digest(row, "operation_id")?), + [] | [_, ..] => return Err(AdminJournalOperationError::Binding), + }; + let before = read_challenge(transaction, operation_id) + .await? + .ok_or(AdminJournalOperationError::Binding)?; + if before.policy_generation != policy_generation { + return Err(AdminJournalOperationError::Conflict); + } + match authorize_challenge( + transaction, + challenge_id, + before.connection_id, + before.operation_id, + policy_generation, + observed_at, + rate_policy, + ) + .await? + { + MycAuthorizationChallengeAuthorization::Resolved(record) + | MycAuthorizationChallengeAuthorization::ExactReplay(record) => Ok(record), + MycAuthorizationChallengeAuthorization::RateLimited => { + Err(AdminJournalOperationError::ResourceExhausted) + } + } +} + +async fn read_pending_decision_operation( + transaction: &mut ServiceSqliteTransaction<'_>, + connection_id: MycConnectionId, + policy_generation: MycConnectionPolicyGeneration, +) -> Result<MycSignerOperationId, ConnectionOperationError> { + let rows = sqlx::query(READ_PENDING_DECISION_FOR_CONNECTION_SQL) + .bind(connection_id.as_bytes().as_slice()) + .bind(policy_generation.sqlite_value()) + .fetch_all(&mut *transaction) + .await + .map_err(|_| ConnectionOperationError::Storage)?; + match rows.as_slice() { + [row] => Ok(MycSignerOperationId::from_persisted(exact_digest( + row, + "operation_id", + )?)), + [] | [_, ..] => Err(ConnectionOperationError::Binding), + } +} + const fn map_governance_error(error: GovernanceOperationError) -> ConnectionOperationError { match error { GovernanceOperationError::Binding => ConnectionOperationError::Binding, diff --git a/src/state_delivery.rs b/src/state_delivery.rs @@ -30,6 +30,30 @@ const ATTEMPT_ID_DOMAIN: &[u8] = b"radroots.myc.delivery_attempt.v1\0"; const READ_ACTIVE_JOB_COUNT_SQL: &str = "SELECT COUNT(*) AS row_count FROM delivery_jobs WHERE status IN ('pending', 'active')"; +const READ_RUNTIME_OUTBOX_STATUS_SQL: &str = r#"SELECT + COUNT(CASE WHEN status IN ('pending', 'active') THEN 1 END) AS pending_count, + (SELECT COUNT(*) FROM delivery_targets WHERE status = 'unknown') AS unknown_count, + MIN(CASE WHEN status IN ('pending', 'active') THEN created_at_unix_ms ELSE NULL END) + AS oldest_pending_at_unix_ms, + typeof(MIN(CASE WHEN status IN ('pending', 'active') THEN created_at_unix_ms ELSE NULL END)) + AS oldest_pending_type +FROM delivery_jobs"#; + +const READ_NEXT_READY_TARGET_SQL: &str = r#"SELECT + CASE WHEN typeof(t.job_id) = 'blob' AND length(t.job_id) = 32 + THEN t.job_id ELSE NULL END AS job_id, + CASE WHEN typeof(t.relay_id) = 'text' + AND length(CAST(t.relay_id AS BLOB)) BETWEEN 1 AND 64 + THEN t.relay_id ELSE NULL END AS relay_id +FROM delivery_targets t +JOIN delivery_jobs j ON j.job_id = t.job_id +WHERE j.status IN ('pending', 'active') + AND t.status IN ('pending', 'retryable', 'unknown') + AND t.active_attempt_id IS NULL + AND (t.next_attempt_at_unix_ms IS NULL OR t.next_attempt_at_unix_ms <= ?) +ORDER BY j.created_at_unix_ms, t.job_id, t.target_index +LIMIT 2"#; + const READ_JOB_SQL: &str = r#"SELECT CASE WHEN typeof(job_id) = 'blob' AND length(job_id) = 32 THEN job_id ELSE NULL END AS job_id, @@ -1000,6 +1024,86 @@ impl fmt::Debug for MycDeliveryClaim { } impl MycStateRepository<'_> { + pub(crate) async fn read_runtime_outbox_status( + &self, + ) -> Result<crate::MycOutboxStatusV1, MycStateRepositoryError> { + let expected = PersistedMetadata::from(self.expected()); + self.host() + .transaction(move |transaction| { + Box::pin(async move { + verify_metadata(transaction, &expected).await?; + let rows = sqlx::query(READ_RUNTIME_OUTBOX_STATUS_SQL) + .fetch_all(&mut *transaction) + .await + .map_err(|_| DeliveryOperationError::Storage)?; + let [row] = rows.as_slice() else { + return Err(DeliveryOperationError::Binding); + }; + let count = |column| { + row.try_get::<i64, _>(column) + .ok() + .and_then(|value| u64::try_from(value).ok()) + .ok_or(DeliveryOperationError::Binding) + }; + let oldest = match row + .try_get::<&str, _>("oldest_pending_type") + .map_err(|_| DeliveryOperationError::Binding)? + { + "null" => None, + "integer" => { + let milliseconds = count("oldest_pending_at_unix_ms")?; + Some( + crate::MycStatusUnixSeconds::new(milliseconds / 1_000) + .map_err(|_| DeliveryOperationError::Binding)?, + ) + } + _ => return Err(DeliveryOperationError::Binding), + }; + Ok(crate::MycOutboxStatusV1::new( + count("pending_count")?, + count("unknown_count")?, + oldest, + )) + }) + }) + .await + .map_err(map_transaction_error) + } + + /// Returns the first exact target eligible for bounded delivery work. + /// + /// Selection is deterministic and performs no claim or network I/O. The + /// subsequent claim transaction remains the sole lease authority. + pub(crate) async fn next_ready_delivery_target( + &self, + observed_at: MycDeliveryTimeUnixMs, + ) -> Result<Option<(MycDeliveryJobId, MycDeliveryRelayId)>, MycStateRepositoryError> { + let expected = PersistedMetadata::from(self.expected()); + self.host() + .transaction(move |transaction| { + Box::pin(async move { + verify_metadata(transaction, &expected).await?; + let rows = sqlx::query(READ_NEXT_READY_TARGET_SQL) + .bind(observed_at.sqlite_value()) + .fetch_all(&mut *transaction) + .await + .map_err(|_| DeliveryOperationError::Storage)?; + match rows.as_slice() { + [] => Ok(None), + [row] => { + let job_id = MycDeliveryJobId::from_persisted(blob32(row, "job_id")?); + let relay_id = MycDeliveryRelayId::new(text(row, "relay_id")?) + .map_err(|_| DeliveryOperationError::Binding)?; + Ok(Some((job_id, relay_id))) + } + _ => Err(DeliveryOperationError::Binding), + } + }) + }) + .await + .map_err(map_transaction_error) + } + /// Claims one eligible target under a bounded expiring attempt lease. pub async fn claim_delivery_target( &self, diff --git a/src/state_discovery.rs b/src/state_discovery.rs @@ -13,6 +13,8 @@ use serde::{Deserialize, Serialize}; use sha2::{Digest, Sha256}; use sqlx::Row; +#[cfg(any(target_os = "linux", target_os = "macos"))] +use crate::state_admin::AdminJournalOperationError; use crate::state_delivery::{ DeliveryOperationError, MycDeliveryArtifactDigest, MycDeliveryJobId, MycDeliveryJobRecord, MycDeliveryJobStatus, MycDeliverySource, MycDeliveryTimeUnixMs, create_job, @@ -24,7 +26,10 @@ use crate::state_repository::{ use crate::{MycExpectedIdentities, MycStateMetadata}; use nostr::RelayUrl as RadrootsNostrRelayUrl; use radroots_nostr::event::{ - Event as RadrootsNostrEvent, Kind as RadrootsNostrKind, Metadata as RadrootsNostrMetadata, + ApplicationHandlerSpec as RadrootsNostrApplicationHandlerSpec, Event as RadrootsNostrEvent, + Kind as RadrootsNostrKind, Metadata as RadrootsNostrMetadata, + Timestamp as RadrootsNostrTimestamp, + build_application_handler as radroots_nostr_build_application_handler_event, }; /// Maximum exact signed-event bytes admitted from the configured event bound. @@ -427,6 +432,52 @@ impl MycDiscoveryPolicies { } } +pub(crate) fn prepare_discovery_signing_bytes( + metadata: &MycStateMetadata, + created_at_unix_seconds: u64, +) -> Result<Box<[u8]>, MycDiscoveryStateError> { + let policy = metadata + .discovery_policies() + .ok_or_else(|| MycDiscoveryStateError::new(MycDiscoveryStateErrorKind::Disabled))?; + let metadata = if policy.metadata_json.is_empty() { + RadrootsNostrMetadata::default() + } else { + serde_json::from_str(&policy.metadata_json).map_err(|_| { + MycDiscoveryStateError::new(MycDiscoveryStateErrorKind::InvalidProjection) + })? + }; + let mut spec = RadrootsNostrApplicationHandlerSpec::new(vec![NIP46_RPC_KIND]) + .with_identifier(policy.handler_identifier.to_string()) + .with_relays( + policy + .public_relays + .iter() + .map(ToString::to_string) + .collect(), + ) + .with_metadata(metadata); + if let Some(url) = &policy.nostrconnect_url { + spec = spec.with_nostr_connect_url(url.to_string()); + } + let public_key = policy + .author_public_key + .parse() + .map_err(|_| MycDiscoveryStateError::new(MycDiscoveryStateErrorKind::IdentityMismatch))?; + let request = radroots_nostr_build_application_handler_event(&spec) + .map_err(|_| MycDiscoveryStateError::new(MycDiscoveryStateErrorKind::InvalidProjection))? + .custom_created_at(RadrootsNostrTimestamp::from_secs(created_at_unix_seconds)) + .into_external_signing_request(public_key) + .map_err(|_| MycDiscoveryStateError::new(MycDiscoveryStateErrorKind::InvalidProjection))?; + let bytes = serde_json::to_vec(&request) + .map_err(|_| MycDiscoveryStateError::new(MycDiscoveryStateErrorKind::InvalidProjection))?; + if bytes.is_empty() || bytes.len() > policy.event_max_bytes { + return Err(MycDiscoveryStateError::new( + MycDiscoveryStateErrorKind::TooLarge, + )); + } + Ok(bytes.into_boxed_slice()) +} + impl fmt::Debug for MycDiscoveryPolicies { fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { formatter @@ -519,6 +570,14 @@ impl MycDiscoveryCommitRequest { self.desired_digest } + pub(crate) const fn event_digest(&self) -> MycDeliveryArtifactDigest { + self.event_digest + } + + pub(crate) const fn projection_digest(&self) -> MycNip05ProjectionDigest { + self.projection_digest + } + fn owned(&self) -> Self { Self { generation_id: self.generation_id, @@ -809,11 +868,33 @@ impl MycStateRepository<'_> { } #[derive(Clone, Copy, Debug, PartialEq, Eq)] -enum DiscoveryOperationError { +pub(crate) enum DiscoveryOperationError { Binding, Storage, } +#[cfg(any(target_os = "linux", target_os = "macos"))] +impl From<DiscoveryOperationError> for AdminJournalOperationError { + fn from(error: DiscoveryOperationError) -> Self { + match error { + DiscoveryOperationError::Binding => Self::Binding, + DiscoveryOperationError::Storage => Self::Storage, + } + } +} + +#[cfg(any(target_os = "linux", target_os = "macos"))] +pub(crate) async fn apply_admin_discovery_publish( + transaction: &mut ServiceSqliteTransaction<'_>, + request: &MycDiscoveryCommitRequest, + delivery_policy: &crate::state_delivery::MycDeliveryPolicies, + discovery_policy: &MycDiscoveryPolicies, +) -> Result<MycDiscoveryCommitAdmission, AdminJournalOperationError> { + commit_desired(transaction, request, delivery_policy, discovery_policy) + .await + .map_err(Into::into) +} + impl From<DeliveryOperationError> for DiscoveryOperationError { fn from(error: DeliveryOperationError) -> Self { match error { diff --git a/src/state_governance.rs b/src/state_governance.rs @@ -342,6 +342,10 @@ impl MycAuditCorrelationId { pub const fn new(bytes: [u8; 32]) -> Self { Self(bytes) } + + pub(crate) const fn as_bytes(&self) -> &[u8; 32] { + &self.0 + } } impl fmt::Debug for MycAuditCorrelationId { @@ -362,7 +366,7 @@ pub enum MycAuditKind { } impl MycAuditKind { - const fn as_str(self) -> &'static str { + pub(crate) const fn as_str(self) -> &'static str { match self { Self::ConnectionAdmission => "connection_admission", Self::ConnectionOperatorDecision => "connection_operator_decision", @@ -373,7 +377,7 @@ impl MycAuditKind { } } - fn parse(value: &str) -> Option<Self> { + pub(crate) fn parse(value: &str) -> Option<Self> { match value { "connection_admission" => Some(Self::ConnectionAdmission), "connection_operator_decision" => Some(Self::ConnectionOperatorDecision), @@ -395,7 +399,7 @@ pub enum MycAuditOutcome { } impl MycAuditOutcome { - const fn as_str(self) -> &'static str { + pub(crate) const fn as_str(self) -> &'static str { match self { Self::Succeeded => "succeeded", Self::Rejected => "rejected", @@ -403,7 +407,7 @@ impl MycAuditOutcome { } } - fn parse(value: &str) -> Option<Self> { + pub(crate) fn parse(value: &str) -> Option<Self> { match value { "succeeded" => Some(Self::Succeeded), "rejected" => Some(Self::Rejected), @@ -430,7 +434,7 @@ pub enum MycAuditReasonCode { } impl MycAuditReasonCode { - const fn as_str(self) -> &'static str { + pub(crate) const fn as_str(self) -> &'static str { match self { Self::Trusted => "trusted", Self::ApprovalRequired => "approval_required", @@ -576,6 +580,47 @@ impl fmt::Debug for MycAuditPage { } } +#[derive(Clone, Copy)] +pub(crate) struct MycAdminAuditQuery { + limit: MycAuditPageLimit, + snapshot_sequence: Option<u64>, + before_sequence: Option<u64>, + from_unix_ms: Option<u64>, + to_unix_ms: Option<u64>, + kind: Option<MycAuditKind>, + outcome: Option<MycAuditOutcome>, +} + +impl MycAdminAuditQuery { + pub(crate) fn new( + limit: MycAuditPageLimit, + snapshot_sequence: Option<u64>, + before_sequence: Option<u64>, + from_unix_ms: Option<u64>, + to_unix_ms: Option<u64>, + kind: Option<MycAuditKind>, + outcome: Option<MycAuditOutcome>, + ) -> Result<Self, MycStateRepositoryError> { + if from_unix_ms.is_some_and(|from| to_unix_ms.is_some_and(|to| from > to)) + || from_unix_ms.is_some_and(|value| i64::try_from(value).is_err()) + || to_unix_ms.is_some_and(|value| i64::try_from(value).is_err()) + { + return Err(MycStateRepositoryError::new( + MycStateRepositoryErrorKind::Binding, + )); + } + Ok(Self { + limit, + snapshot_sequence, + before_sequence, + from_unix_ms, + to_unix_ms, + kind, + outcome, + }) + } +} + /// Explicit bounded retention and compaction policy. #[derive(Clone, Copy, Debug, PartialEq, Eq)] pub struct MycGovernanceCompactionPolicy { @@ -766,6 +811,82 @@ pub(crate) async fn record_audit( } impl MycStateRepository<'_> { + pub(crate) async fn read_admin_audit_page( + &self, + query: MycAdminAuditQuery, + ) -> Result<MycAuditPage, MycStateRepositoryError> { + let expected = PersistedMetadata::from(self.expected()); + self.host() + .transaction(move |transaction| { + Box::pin(async move { + verify_metadata(transaction, &expected).await?; + read_admin_audit_page(transaction, query).await + }) + }) + .await + .map_err(map_transaction_error) + } + + pub(crate) async fn read_admin_audit_summary( + &self, + from_unix_ms: u64, + to_unix_ms: u64, + ) -> Result<Box<[(String, u64)]>, MycStateRepositoryError> { + if from_unix_ms > to_unix_ms + || i64::try_from(from_unix_ms).is_err() + || i64::try_from(to_unix_ms).is_err() + { + return Err(MycStateRepositoryError::new( + MycStateRepositoryErrorKind::Binding, + )); + } + let expected = PersistedMetadata::from(self.expected()); + self.host() + .transaction(move |transaction| { + Box::pin(async move { + verify_metadata(transaction, &expected).await?; + let rows = sqlx::query( + r#"SELECT audit_kind, outcome, COUNT(*) AS item_count + FROM operation_audit + WHERE occurred_at_unix_ms BETWEEN ? AND ? + GROUP BY audit_kind, outcome + ORDER BY audit_kind ASC, outcome ASC + LIMIT 19"#, + ) + .bind( + i64::try_from(from_unix_ms) + .map_err(|_| GovernanceOperationError::Binding)?, + ) + .bind(i64::try_from(to_unix_ms).map_err(|_| GovernanceOperationError::Binding)?) + .fetch_all(&mut *transaction) + .await + .map_err(|_| GovernanceOperationError::Storage)?; + let mut counts = Vec::with_capacity(rows.len()); + for row in rows { + let kind = row + .try_get::<&str, _>("audit_kind") + .ok() + .and_then(MycAuditKind::parse) + .ok_or(GovernanceOperationError::Binding)?; + let outcome = row + .try_get::<&str, _>("outcome") + .ok() + .and_then(MycAuditOutcome::parse) + .ok_or(GovernanceOperationError::Binding)?; + let count = row + .try_get::<i64, _>("item_count") + .ok() + .and_then(|value| u64::try_from(value).ok()) + .ok_or(GovernanceOperationError::Binding)?; + counts.push((format!("{}.{}", kind.as_str(), outcome.as_str()), count)); + } + Ok(counts.into_boxed_slice()) + }) + }) + .await + .map_err(map_transaction_error) + } + /// Reads one deterministic, snapshot-bounded page of safe audit evidence. pub async fn read_audit_page( &self, @@ -1012,6 +1133,82 @@ async fn read_audit_page( }) } +#[allow(clippy::too_many_arguments)] +async fn read_admin_audit_page( + transaction: &mut ServiceSqliteTransaction<'_>, + query: MycAdminAuditQuery, +) -> Result<MycAuditPage, GovernanceOperationError> { + let high_water = read_audit_state(transaction).await?; + let snapshot = query.snapshot_sequence.unwrap_or(high_water); + if snapshot > high_water + || query + .before_sequence + .is_some_and(|value| value == 0 || value > snapshot.saturating_add(1)) + { + return Err(GovernanceOperationError::Binding); + } + let before = query + .before_sequence + .unwrap_or_else(|| snapshot.saturating_add(1)); + let fetch_limit = u32::from(query.limit.0) + 1; + let from = query + .from_unix_ms + .map(i64::try_from) + .transpose() + .map_err(|_| GovernanceOperationError::Binding)?; + let to = query + .to_unix_ms + .map(i64::try_from) + .transpose() + .map_err(|_| GovernanceOperationError::Binding)?; + let kind = query.kind.map(MycAuditKind::as_str); + let outcome = query.outcome.map(MycAuditOutcome::as_str); + let rows = sqlx::query( + r#"SELECT audit_id FROM operation_audit + WHERE audit_sequence <= ? AND audit_sequence < ? + AND (? IS NULL OR occurred_at_unix_ms >= ?) + AND (? IS NULL OR occurred_at_unix_ms <= ?) + AND (? IS NULL OR audit_kind = ?) + AND (? IS NULL OR outcome = ?) + ORDER BY audit_sequence DESC LIMIT ?"#, + ) + .bind(to_i64(snapshot)?) + .bind(to_i64(before)?) + .bind(from) + .bind(from) + .bind(to) + .bind(to) + .bind(kind) + .bind(kind) + .bind(outcome) + .bind(outcome) + .bind(i64::from(fetch_limit)) + .fetch_all(&mut *transaction) + .await + .map_err(|_| GovernanceOperationError::Storage)?; + if rows.len() > usize::try_from(fetch_limit).map_err(|_| GovernanceOperationError::Binding)? { + return Err(GovernanceOperationError::Binding); + } + let has_more = rows.len() > usize::from(query.limit.0); + let mut items = Vec::with_capacity(rows.len().min(usize::from(query.limit.0))); + for row in rows.iter().take(usize::from(query.limit.0)) { + let id = bounded_digest(row, "audit_id")?; + items.push( + read_audit_by_id(transaction, id) + .await? + .ok_or(GovernanceOperationError::Binding)?, + ); + } + let next = has_more + .then(|| items.last().map(|item| item.sequence)) + .flatten(); + Ok(MycAuditPage { + snapshot_sequence: snapshot, + items: items.into_boxed_slice(), + next_before_sequence: next, + }) +} + async fn compact_governance( transaction: &mut ServiceSqliteTransaction<'_>, observed_at: MycConnectionTimeUnixMs, diff --git a/src/state_request.rs b/src/state_request.rs @@ -453,10 +453,6 @@ impl MycSignerRequest { self.request_digest } - pub(crate) const fn received_at(&self) -> MycRequestReceivedAtUnixMs { - self.received_at - } - fn derived_operation_id(&self) -> MycSignerOperationId { MycSignerOperationId(derive_operation_id( &self.request_identity, @@ -516,6 +512,10 @@ impl MycSignerRequestRecord { &self.client_public_key } + pub(crate) const fn request_id(&self) -> &MycNip46RequestId { + &self.request_id + } + #[must_use] /// Returns the stable logical operation identity. pub const fn operation_id(&self) -> MycSignerOperationId { @@ -554,7 +554,7 @@ impl MycSignerRequestRecord { && self.first_event_id == request.event_id() && self.method == request.method() && self.request_digest == request.request_digest() - && self.received_at == request.received_at() + && self.received_at == request.received_at } pub(crate) const fn received_at(&self) -> MycRequestReceivedAtUnixMs { diff --git a/src/state_response.rs b/src/state_response.rs @@ -463,6 +463,28 @@ impl MycStateRepository<'_> { .await .map_err(map_transaction_error) } + + /// Reads an already committed response by stable signer operation identity. + /// + /// Runtime replay handling uses this lookup before any provider call so an + /// exact completed replay always reuses the originally committed bytes. + pub(crate) async fn read_nip46_response_by_operation( + &self, + operation_id: MycSignerOperationId, + ) -> Result<Option<MycNip46ResponseRecord>, MycStateRepositoryError> { + let expected = PersistedMetadata::from(self.expected()); + self.host() + .transaction(move |transaction| { + Box::pin(async move { + require_expected_metadata(transaction, &expected) + .await + .map_err(AtomicOperationError::from)?; + read_response_by_operation(transaction, operation_id).await + }) + }) + .await + .map_err(map_transaction_error) + } } #[derive(Clone, Copy, Debug, PartialEq, Eq)] diff --git a/src/transport_nostr_adapter.rs b/src/transport_nostr_adapter.rs @@ -1,27 +1,28 @@ //! Exact source-locked Nostr delivery adapter owned by the Myc runtime. -#![allow( - dead_code, - reason = "Step 159 Unit 12 seals the adapter before Unit 15 runtime graph wiring" -)] - use core::{fmt, future::Future, pin::Pin}; use std::{collections::BTreeMap, error::Error}; use radroots_event_codec::Codec; use radroots_transport::{ - EventSource, FetchRequest, Target, TargetSet, + EventSource, EventSubscriber, FetchRequest, Target, TargetSet, outcome::{DeliveryOutcomeKind, FetchTargetState}, policy::{SatisfactionClass, SatisfactionPolicy, TargetPolicy}, sink::{DeliveryPayload, DeliveryRequest, DeliveryTargetReceipt}, - source::{FetchBounds, FetchSelector}, + source::{ + BoxSubscription, FetchBounds, FetchSelector, SubscriptionBounds, SubscriptionCheckpoint, + SubscriptionRequest, + }, + target::TargetFingerprint, }; use radroots_transport_nostr::{ Config, NostrTransport, PreparedDelivery, RelayAccess, RelayEndpoint, RelayProfile, RelayProfileKind, RelayUrlPolicy, }; -use crate::{MycConfigDocumentV1, MycConfigProfile, MycDeliveryRelayId}; +use crate::{MycConfigDocumentV1, MycConfigProfile, MycDeliveryRelayId, MycRateRelayId}; + +const NIP46_RPC_KIND: u32 = 24_133; pub(crate) type RelayExecutionFuture<'a> = Pin< Box<dyn Future<Output = Result<MycRelayExecutionOutcome, MycRelayAdapterError>> + Send + 'a>, @@ -41,6 +42,7 @@ pub(crate) struct MycRelayAdapterError { } impl MycRelayAdapterError { + #[cfg(test)] pub(crate) const fn kind(&self) -> MycRelayAdapterErrorKind { self.kind } @@ -67,6 +69,12 @@ const fn adapter_error(kind: MycRelayAdapterErrorKind) -> MycRelayAdapterError { MycRelayAdapterError { kind } } +pub(crate) const fn runtime_relay_adapter_error( + kind: MycRelayAdapterErrorKind, +) -> MycRelayAdapterError { + adapter_error(kind) +} + #[derive(Clone, Copy, Debug, PartialEq, Eq)] pub(crate) enum MycRelayExecutionOutcome { Accepted, @@ -109,6 +117,208 @@ struct RelayTarget { target: Target, } +/// One exact source-locked subscription group owned by the sole ingress task. +pub(crate) struct MycNostrIngressAdapter { + groups: Box<[IngressGroup]>, +} + +struct IngressGroup { + transport: NostrTransport, + targets: TargetSet, + relay_ids: BTreeMap<TargetFingerprint, MycRateRelayId>, + required: bool, + request_id: &'static str, +} + +impl MycNostrIngressAdapter { + pub(crate) fn from_configuration( + configuration: &MycConfigDocumentV1, + ) -> Result<Self, MycRelayAdapterError> { + let relays = configuration + .normalized() + .pointer("/relays") + .and_then(serde_json::Value::as_array) + .ok_or_else(|| adapter_error(MycRelayAdapterErrorKind::Configuration))?; + let connect_timeout = + configuration_integer(configuration, "/transport/connect_deadline_ms")?; + let request_timeout = + configuration_integer(configuration, "/transport/ingress/subscription_deadline_ms")?; + let mut public = Vec::new(); + let mut public_targets = Vec::new(); + let mut public_ids = BTreeMap::new(); + let mut public_required = false; + let mut local = Vec::new(); + let mut local_targets = Vec::new(); + let mut local_ids = BTreeMap::new(); + let mut local_required = false; + for relay in relays.iter().filter(|relay| { + relay.pointer("/read").and_then(serde_json::Value::as_bool) == Some(true) + }) { + let id = relay + .pointer("/id") + .and_then(serde_json::Value::as_str) + .and_then(|value| MycRateRelayId::new(value).ok()) + .ok_or_else(|| adapter_error(MycRelayAdapterErrorKind::Configuration))?; + let url = relay + .pointer("/url") + .and_then(serde_json::Value::as_str) + .ok_or_else(|| adapter_error(MycRelayAdapterErrorKind::Configuration))?; + let required = relay + .pointer("/required") + .and_then(serde_json::Value::as_bool) + .ok_or_else(|| adapter_error(MycRelayAdapterErrorKind::Configuration))?; + let target = Target::nostr_relay(url) + .map_err(|_| adapter_error(MycRelayAdapterErrorKind::Target))?; + let fingerprint = target.fingerprint().clone(); + let (kind, policy) = if url.starts_with("wss://") { + (RelayProfileKind::Public, RelayUrlPolicy::Public) + } else if configuration.profile() == MycConfigProfile::RepoLocal + && url.starts_with("ws://") + { + (RelayProfileKind::Simulator, RelayUrlPolicy::Local) + } else { + return Err(adapter_error(MycRelayAdapterErrorKind::Configuration)); + }; + let endpoint = RelayEndpoint::new(url, policy, RelayAccess::ReadOnly) + .map_err(|_| adapter_error(MycRelayAdapterErrorKind::Configuration))?; + match kind { + RelayProfileKind::Public => { + public.push(endpoint); + public_targets.push(target); + public_required |= required; + if public_ids.insert(fingerprint, id).is_some() { + return Err(adapter_error(MycRelayAdapterErrorKind::Configuration)); + } + } + RelayProfileKind::Simulator => { + local.push(endpoint); + local_targets.push(target); + local_required |= required; + if local_ids.insert(fingerprint, id).is_some() { + return Err(adapter_error(MycRelayAdapterErrorKind::Configuration)); + } + } + RelayProfileKind::Device => { + return Err(adapter_error(MycRelayAdapterErrorKind::Configuration)); + } + _ => return Err(adapter_error(MycRelayAdapterErrorKind::Configuration)), + } + } + let mut groups = Vec::with_capacity(2); + add_ingress_group( + &mut groups, + RelayProfileKind::Public, + public, + public_targets, + public_ids, + public_required, + connect_timeout, + request_timeout, + "myc-runtime-public", + )?; + add_ingress_group( + &mut groups, + RelayProfileKind::Simulator, + local, + local_targets, + local_ids, + local_required, + connect_timeout, + request_timeout, + "myc-runtime-local", + )?; + if groups.is_empty() || groups.len() > 2 || !groups.iter().any(|group| group.required) { + return Err(adapter_error(MycRelayAdapterErrorKind::Configuration)); + } + Ok(Self { + groups: groups.into_boxed_slice(), + }) + } + + pub(crate) fn group_count(&self) -> usize { + self.groups.len() + } + + pub(crate) fn group_is_required(&self, index: usize) -> Option<bool> { + self.groups.get(index).map(|group| group.required) + } + + pub(crate) async fn subscribe( + &self, + index: usize, + event_limit: u16, + deadline_unix_ms: u64, + checkpoints: &[SubscriptionCheckpoint], + ) -> Result<BoxSubscription, MycRelayAdapterError> { + let group = self + .groups + .get(index) + .ok_or_else(|| adapter_error(MycRelayAdapterErrorKind::Configuration))?; + let selector = FetchSelector::all() + .with_kinds(vec![NIP46_RPC_KIND]) + .map_err(|_| adapter_error(MycRelayAdapterErrorKind::Configuration))?; + let request = SubscriptionRequest::new( + group.request_id, + group.targets.clone(), + SubscriptionBounds::new(event_limit, deadline_unix_ms) + .map_err(|_| adapter_error(MycRelayAdapterErrorKind::Configuration))?, + ) + .map_err(|_| adapter_error(MycRelayAdapterErrorKind::Configuration))? + .with_selector(selector) + .with_checkpoints(checkpoints.iter().cloned()) + .map_err(|_| adapter_error(MycRelayAdapterErrorKind::Configuration))?; + group + .transport + .subscribe(request) + .await + .map_err(|_| adapter_error(MycRelayAdapterErrorKind::Execution)) + } + + pub(crate) fn relay_id( + &self, + index: usize, + target: &TargetFingerprint, + ) -> Result<MycRateRelayId, MycRelayAdapterError> { + self.groups + .get(index) + .and_then(|group| group.relay_ids.get(target)) + .cloned() + .ok_or_else(|| adapter_error(MycRelayAdapterErrorKind::Target)) + } +} + +#[allow(clippy::too_many_arguments)] +fn add_ingress_group( + groups: &mut Vec<IngressGroup>, + kind: RelayProfileKind, + endpoints: Vec<RelayEndpoint>, + targets: Vec<Target>, + relay_ids: BTreeMap<TargetFingerprint, MycRateRelayId>, + required: bool, + connect_timeout: u64, + request_timeout: u64, + request_id: &'static str, +) -> Result<(), MycRelayAdapterError> { + if targets.is_empty() { + return Ok(()); + } + let transport = build_transport(kind, endpoints, connect_timeout, request_timeout)? + .ok_or_else(|| adapter_error(MycRelayAdapterErrorKind::Configuration))?; + let targets = + TargetSet::new(targets).map_err(|_| adapter_error(MycRelayAdapterErrorKind::Target))?; + if targets.len() != relay_ids.len() { + return Err(adapter_error(MycRelayAdapterErrorKind::Configuration)); + } + groups.push(IngressGroup { + transport, + targets, + relay_ids, + required, + request_id, + }); + Ok(()) +} + impl MycNostrDeliveryAdapter { pub(crate) fn from_configuration( configuration: &MycConfigDocumentV1, diff --git a/tests/build_policy.rs b/tests/build_policy.rs @@ -23,7 +23,7 @@ fn service_host_is_the_exact_default_feature_profile() { assert!(MANIFEST.contains("[features]\ndefault = [\"service-host\"]\nservice-host = []")); assert!(!MANIFEST.contains("getrandom = \"0.2\"")); assert!(MANIFEST.contains( - "tokio = { version = \"1.48\", default-features = false, features = [\"io-util\", \"macros\", \"net\", \"rt-multi-thread\", \"sync\", \"time\"] }" + "tokio = { version = \"1.48\", default-features = false, features = [\"io-util\", \"macros\", \"net\", \"rt-multi-thread\", \"signal\", \"sync\", \"time\"] }" )); } @@ -35,6 +35,13 @@ fn shared_runtime_paths_is_exactly_pinned_to_the_source_locked_lib() { } #[test] +fn shared_identity_is_exactly_pinned_to_the_source_locked_lib() { + assert!(MANIFEST.contains( + "radroots_identity = { git = \"https://github.com/radrootslabs/lib\", rev = \"7d7b454b4c9ed86569671993bd03ca868b676665\", version = \"=0.1.0-alpha\" }" + )); +} + +#[test] fn shared_service_host_is_exactly_pinned_to_the_source_locked_lib() { assert!(MANIFEST.contains( "radroots_service_host = { git = \"https://github.com/radrootslabs/lib\", rev = \"7d7b454b4c9ed86569671993bd03ca868b676665\", version = \"=0.1.0-alpha\" }" diff --git a/tests/package_boundary.rs b/tests/package_boundary.rs @@ -34,6 +34,9 @@ const PROCESS_V1_UNSUPPORTED: &str = include_str!("../src/process_v1_unsupported const CONFIG_LOADER: &str = include_str!("../src/config_loader.rs"); const SYSTEM_DOCTOR: &str = include_str!("../src/system_doctor.rs"); const STATUS_V1: &str = include_str!("../src/status_v1.rs"); +const RUNTIME_GRAPH: &str = include_str!("../src/runtime_graph.rs"); +const RUNTIME_NIP46: &str = include_str!("../src/runtime_nip46.rs"); +const RUNTIME_SIGNAL: &str = include_str!("../src/runtime_signal.rs"); const RUNTIME_SUPERVISION: &str = include_str!("../src/runtime_supervision.rs"); const RUNTIME_SUPERVISION_CONTRACT: &str = include_str!("../contracts/services_hardening/runtime_supervision.v1.json"); @@ -84,6 +87,9 @@ const SOURCES: &[&str] = &[ include_str!("../src/provider_verification.rs"), include_str!("../src/runtime_context.rs"), include_str!("../src/runtime_foundation.rs"), + include_str!("../src/runtime_graph.rs"), + include_str!("../src/runtime_nip46.rs"), + include_str!("../src/runtime_signal.rs"), include_str!("../src/runtime_supervision.rs"), include_str!("../src/status_v1.rs"), include_str!("../src/system_doctor.rs"), @@ -139,6 +145,9 @@ fn implementation_modules_are_private_and_rustdoc_uses_the_reviewed_readme() { "provider_verification", "runtime_context", "runtime_foundation", + "runtime_graph", + "runtime_nip46", + "runtime_signal", "runtime_supervision", "status_v1", "system_doctor", @@ -185,13 +194,49 @@ fn implementation_modules_are_private_and_rustdoc_uses_the_reviewed_readme() { "Restore derives expected backup identity\nfrom the trusted manifest digest", "The production adapter composes secure path and disk inspection", "It never publishes a relay event", - "runtime task-graph wiring and startup\nhandshakes remain the next ordered Step 159 unit", + "The production `run` path owns the exact five-role bounded graph", + "Required relay subscriptions and provider handshakes complete before\nReady", + "one configured absolute graceful-shutdown deadline", ] { assert!(README.contains(required), "README is missing `{required}`"); } } #[test] +fn step159_runtime_graph_is_fixed_joined_and_binary_signal_owned() { + for required in [ + "const TASK_ADMIN_SERVER: &str = \"admin_server\"", + "const TASK_OPERATIONS_SERVER: &str = \"operations_server\"", + "const TASK_RELAY_INGRESS: &str = \"relay_ingress\"", + "const TASK_PROVIDER_DISPATCH: &str = \"provider_dispatch\"", + "const TASK_DELIVERY_OUTBOX: &str = \"delivery_outbox\"", + "open_initial_subscriptions(&ingress_adapter, &configuration).await", + "required_relays_ready(slots)", + "MycRuntimeNip46Coordinator::new", + "let admission_evidence = runtime_nip46_admission_evidence()", + "let mut retry = initial", + "GracefulShutdown::new(grace)", + "ProcessSignalAdapter::new(HostSignalSource::new(signals))", + "publish(MycServicePhase::Ready)", + ] { + assert!(RUNTIME_GRAPH.contains(required), "missing `{required}`"); + } + assert_eq!(RUNTIME_GRAPH.matches(".spawn(").count(), 5); + assert!(RUNTIME_NIP46.contains("commit_nip46_response(&commit)")); + assert!(RUNTIME_NIP46.contains("ExactCompletedReplay")); + assert!(RUNTIME_NIP46.contains("admission_evidence: MycRuntimeNip46AdmissionEvidence")); + assert!(RUNTIME_SIGNAL.contains("pub trait MycProcessSignalSource: Send")); + for forbidden in [ + "tokio::spawn", + "spawn_blocking", + "std::thread::spawn", + "std::process::exit", + ] { + assert!(!RUNTIME_GRAPH.contains(forbidden), "found `{forbidden}`"); + } +} + +#[test] fn reviewed_api_is_root_only_and_exposes_no_implementation_authority() { for required in [ "pub struct myc::MycCliExecutionPlanV1", @@ -709,7 +754,7 @@ fn step158_runtime_supervision_is_one_owned_bounded_redacted_graph() { for required in [ "one sealed, bounded critical-task graph", "task names and handles remain internal", - "Step 159 owns\nthe process panic hook, signal installation", + "The binary owns signal\ninstallation", ] { assert!(README.contains(required), "README is missing `{required}`"); } diff --git a/tests/services_hardening_config_contract.rs b/tests/services_hardening_config_contract.rs @@ -181,6 +181,13 @@ fn semantic_valid(value: &Value, profile: Profile) -> bool { { return false; } + let ingress = &value["transport"]["ingress"]; + if !(1000..=300000).contains(&ingress["subscription_deadline_ms"].as_u64().unwrap_or(0)) + || ingress["maximum_past_seconds"].as_u64() > Some(3600) + || ingress["maximum_future_seconds"].as_u64() > Some(3600) + { + return false; + } let required_writers = relays .iter() .filter(|relay| { @@ -718,6 +725,23 @@ fn bounds_relationships_and_conditional_authority_fail_closed() { backoff["transport"]["publish_retry"]["initial_backoff_ms"] = json!(30_000); backoff["transport"]["publish_retry"]["maximum_backoff_ms"] = json!(1); assert_rejected(&backoff, Profile::Production); + for (field, invalid_values) in [ + ("subscription_deadline_ms", vec![json!(999), json!(300001)]), + ("maximum_past_seconds", vec![json!(3601), json!(-1)]), + ("maximum_future_seconds", vec![json!(3601), json!(-1)]), + ] { + for invalid_value in invalid_values { + let mut invalid = value.clone(); + invalid["transport"]["ingress"][field] = invalid_value; + assert_rejected(&invalid, Profile::Production); + } + let mut missing = value.clone(); + missing["transport"]["ingress"] + .as_object_mut() + .expect("ingress object") + .remove(field); + assert_rejected(&missing, Profile::Production); + } let mut quorum = value.clone(); quorum["transport"]["delivery_policy"] = json!({"mode": "required_quorum", "required_acknowledgements": 3}); @@ -783,6 +807,9 @@ fn schema_and_version_are_closed_and_defaults_are_only_safe_leaves() { "challenges", "rate_limits", "delivery_policy", + "subscription_deadline_ms", + "maximum_past_seconds", + "maximum_future_seconds", "discovery", ] { assert!( diff --git a/tests/services_hardening_diagnostics.rs b/tests/services_hardening_diagnostics.rs @@ -129,20 +129,20 @@ fn binary_writes_only_fixed_json_diagnostics_to_stderr() { assert!(!invalid_stderr.contains(canary)); let repo_local = tempfile::tempdir().expect("repo-local root"); - let unavailable = Command::new(env!("CARGO_BIN_EXE_myc")) + let unconfigured = Command::new(env!("CARGO_BIN_EXE_myc")) .args(["--profile", "repo-local", "--instance", "primary"]) .arg("--repo-local-root") .arg(repo_local.path()) .arg("run") .output() .expect("admitted invocation"); - assert_eq!(unavailable.status.code(), Some(3)); - assert!(unavailable.stdout.is_empty()); + assert_eq!(unconfigured.status.code(), Some(2)); + assert!(unconfigured.stdout.is_empty()); assert_eq!( - String::from_utf8(unavailable.stderr).expect("unavailable stderr"), + String::from_utf8(unconfigured.stderr).expect("unconfigured stderr"), format!( "{}\n", - MycLogRecord::process_result(MycProcessResult::ServiceOrDependencyUnavailable) + MycLogRecord::process_result(MycProcessResult::InputOrConfiguration) ) ); } diff --git a/tests/services_hardening_legacy_removal.rs b/tests/services_hardening_legacy_removal.rs @@ -184,7 +184,6 @@ fn obsolete_provider_sources_and_dependencies_are_absent() { "axum", "keyring", "nostr-sdk", - "radroots_identity", "radroots_event", "radroots_signing", "rand", @@ -246,10 +245,10 @@ fn binary_uses_only_the_hardened_parser_and_fails_closed_before_dispatch() { .args(["--profile", "service-host", "--instance", "primary", "run"]) .output() .expect("run admitted command"); - assert_eq!(admitted.status.code(), Some(3)); + assert_eq!(admitted.status.code(), Some(2)); assert_eq!( String::from_utf8(admitted.stderr).expect("utf8 stderr"), - process_diagnostic(MycProcessResult::ServiceOrDependencyUnavailable) + process_diagnostic(MycProcessResult::InputOrConfiguration) ); } diff --git a/tests/services_hardening_native_release.rs b/tests/services_hardening_native_release.rs @@ -149,7 +149,7 @@ fn every_radroots_dependency_is_exactly_source_locked() { .iter() .filter(|(name, _)| name.starts_with("radroots_")) .collect::<Vec<_>>(); - assert_eq!(radroots.len(), 10); + assert_eq!(radroots.len(), 11); for (name, dependency) in radroots { let dependency = dependency.as_table().expect("detailed dependency"); assert_eq!( diff --git a/tests/services_hardening_process.rs b/tests/services_hardening_process.rs @@ -142,6 +142,24 @@ fn binary_executes_config_state_backup_restore_and_doctor_boundaries() { assert_eq!(status_value["schema_version"], 11); assert_eq!(status_value["integrity"], "verified"); + let service_status = fixture.run(&["status"]); + assert_success(&service_status); + let service_status_value: serde_json::Value = + serde_json::from_slice(&service_status.stdout).expect("service status JSON"); + assert_eq!( + service_status_value["myc"]["connection_counts"], + serde_json::json!({ + "active": 0, + "denied": 0, + "expired": 0, + "pending": 0, + }) + ); + assert_eq!( + service_status_value["myc"]["outbox"], + serde_json::json!({"pending": 0, "unknown": 0}) + ); + let bundle = fixture.root.path().join("backup"); let backup = fixture .command(&["state", "backup", "--operation-id", "process-backup-01"]) diff --git a/tests/services_hardening_state_metadata.rs b/tests/services_hardening_state_metadata.rs @@ -89,7 +89,7 @@ fn exact_database_configuration_identity_and_policy_bindings_are_frozen() { assert_eq!(versions.status(), MYC_SIGNER_STATUS_CONTRACT_VERSION); assert_eq!( hex::encode(metadata.configuration_digest().as_bytes()), - "56942af2ea11124114cae734dacdac75efd22c5476af6fe970e1d586d439b630" + "68f65b32652a4646d33aafc12075ff44548fb7d9f9a96c4eacaa823702d386cc" ); } @@ -123,6 +123,20 @@ fn digest_uses_fully_defaulted_values_and_changes_with_normalized_policy() { explicit.configuration_digest(), changed.configuration_digest() ); + + let changed_ingress = state_metadata( + &runtime, + &EXAMPLE.replacen( + "maximum_past_seconds = 120", + "maximum_past_seconds = 121", + 1, + ), + ) + .expect("changed ingress policy"); + assert_ne!( + explicit.configuration_digest(), + changed_ingress.configuration_digest() + ); } #[test]