transport_nostr_adapter.rs (27486B)
1 //! Exact source-locked Nostr delivery adapter owned by the Myc runtime. 2 3 use core::{fmt, future::Future, pin::Pin}; 4 use std::{collections::BTreeMap, error::Error}; 5 6 use radroots_event_codec::Codec; 7 use radroots_transport::{ 8 EventSource, EventSubscriber, FetchRequest, Target, TargetSet, 9 outcome::{DeliveryOutcomeKind, FetchTargetState}, 10 policy::{SatisfactionClass, SatisfactionPolicy, TargetPolicy}, 11 sink::{DeliveryPayload, DeliveryRequest, DeliveryTargetReceipt}, 12 source::{ 13 BoxSubscription, FetchBounds, FetchSelector, SubscriptionBounds, SubscriptionCheckpoint, 14 SubscriptionRequest, 15 }, 16 target::TargetFingerprint, 17 }; 18 use radroots_transport_nostr::{ 19 Config, NostrTransport, PreparedDelivery, RelayAccess, RelayEndpoint, RelayProfile, 20 RelayProfileKind, RelayUrlPolicy, 21 }; 22 23 use crate::{MycConfigDocumentV1, MycConfigProfile, MycDeliveryRelayId, MycRateRelayId}; 24 25 const NIP46_RPC_KIND: u32 = 24_133; 26 27 pub(crate) type RelayExecutionFuture<'a> = Pin< 28 Box<dyn Future<Output = Result<MycRelayExecutionOutcome, MycRelayAdapterError>> + Send + 'a>, 29 >; 30 31 #[derive(Clone, Copy, Debug, PartialEq, Eq)] 32 pub(crate) enum MycRelayAdapterErrorKind { 33 Configuration, 34 Target, 35 Payload, 36 Preparation, 37 Execution, 38 } 39 40 pub(crate) struct MycRelayAdapterError { 41 kind: MycRelayAdapterErrorKind, 42 } 43 44 impl MycRelayAdapterError { 45 #[cfg(test)] 46 pub(crate) const fn kind(&self) -> MycRelayAdapterErrorKind { 47 self.kind 48 } 49 } 50 51 impl fmt::Debug for MycRelayAdapterError { 52 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 53 formatter 54 .debug_struct("MycRelayAdapterError") 55 .field("kind", &self.kind) 56 .finish() 57 } 58 } 59 60 impl fmt::Display for MycRelayAdapterError { 61 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 62 formatter.write_str("Myc relay adapter failed") 63 } 64 } 65 66 impl Error for MycRelayAdapterError {} 67 68 const fn adapter_error(kind: MycRelayAdapterErrorKind) -> MycRelayAdapterError { 69 MycRelayAdapterError { kind } 70 } 71 72 pub(crate) const fn runtime_relay_adapter_error( 73 kind: MycRelayAdapterErrorKind, 74 ) -> MycRelayAdapterError { 75 adapter_error(kind) 76 } 77 78 #[derive(Clone, Copy, Debug, PartialEq, Eq)] 79 pub(crate) enum MycRelayExecutionOutcome { 80 Accepted, 81 Rejected, 82 TransportFailed, 83 UnknownAcknowledgement, 84 } 85 86 pub(crate) struct MycPreparedRelayDelivery { 87 transport: NostrTransport, 88 prepared: PreparedDelivery, 89 } 90 91 impl fmt::Debug for MycPreparedRelayDelivery { 92 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 93 formatter.write_str("MycPreparedRelayDelivery([redacted])") 94 } 95 } 96 97 pub(crate) trait MycRelayAdapter: Send + Sync { 98 type Prepared: Send; 99 100 fn prepare( 101 &self, 102 relay_id: &MycDeliveryRelayId, 103 request_id: String, 104 exact_event_bytes: &[u8], 105 deadline_unix_ms: u64, 106 ) -> Result<Self::Prepared, MycRelayAdapterError>; 107 108 fn execute<'a>(&'a self, prepared: Self::Prepared) -> RelayExecutionFuture<'a>; 109 } 110 111 pub(crate) struct MycNostrDeliveryAdapter { 112 targets: BTreeMap<MycDeliveryRelayId, RelayTarget>, 113 } 114 115 struct RelayTarget { 116 transport: NostrTransport, 117 target: Target, 118 } 119 120 /// One exact source-locked subscription group owned by the sole ingress task. 121 pub(crate) struct MycNostrIngressAdapter { 122 groups: Box<[IngressGroup]>, 123 } 124 125 struct IngressGroup { 126 transport: NostrTransport, 127 targets: TargetSet, 128 relay_ids: BTreeMap<TargetFingerprint, MycRateRelayId>, 129 required: bool, 130 request_id: &'static str, 131 } 132 133 impl MycNostrIngressAdapter { 134 pub(crate) fn from_configuration( 135 configuration: &MycConfigDocumentV1, 136 ) -> Result<Self, MycRelayAdapterError> { 137 let relays = configuration 138 .normalized() 139 .pointer("/relays") 140 .and_then(serde_json::Value::as_array) 141 .ok_or_else(|| adapter_error(MycRelayAdapterErrorKind::Configuration))?; 142 let connect_timeout = 143 configuration_integer(configuration, "/transport/connect_deadline_ms")?; 144 let request_timeout = 145 configuration_integer(configuration, "/transport/ingress/subscription_deadline_ms")?; 146 let mut public = Vec::new(); 147 let mut public_targets = Vec::new(); 148 let mut public_ids = BTreeMap::new(); 149 let mut public_required = false; 150 let mut local = Vec::new(); 151 let mut local_targets = Vec::new(); 152 let mut local_ids = BTreeMap::new(); 153 let mut local_required = false; 154 for relay in relays.iter().filter(|relay| { 155 relay.pointer("/read").and_then(serde_json::Value::as_bool) == Some(true) 156 }) { 157 let id = relay 158 .pointer("/id") 159 .and_then(serde_json::Value::as_str) 160 .and_then(|value| MycRateRelayId::new(value).ok()) 161 .ok_or_else(|| adapter_error(MycRelayAdapterErrorKind::Configuration))?; 162 let url = relay 163 .pointer("/url") 164 .and_then(serde_json::Value::as_str) 165 .ok_or_else(|| adapter_error(MycRelayAdapterErrorKind::Configuration))?; 166 let required = relay 167 .pointer("/required") 168 .and_then(serde_json::Value::as_bool) 169 .ok_or_else(|| adapter_error(MycRelayAdapterErrorKind::Configuration))?; 170 let target = Target::nostr_relay(url) 171 .map_err(|_| adapter_error(MycRelayAdapterErrorKind::Target))?; 172 let fingerprint = target.fingerprint().clone(); 173 let (kind, policy) = if url.starts_with("wss://") { 174 (RelayProfileKind::Public, RelayUrlPolicy::Public) 175 } else if configuration.profile() == MycConfigProfile::RepoLocal 176 && url.starts_with("ws://") 177 { 178 (RelayProfileKind::Simulator, RelayUrlPolicy::Local) 179 } else { 180 return Err(adapter_error(MycRelayAdapterErrorKind::Configuration)); 181 }; 182 let endpoint = RelayEndpoint::new(url, policy, RelayAccess::ReadOnly) 183 .map_err(|_| adapter_error(MycRelayAdapterErrorKind::Configuration))?; 184 match kind { 185 RelayProfileKind::Public => { 186 public.push(endpoint); 187 public_targets.push(target); 188 public_required |= required; 189 if public_ids.insert(fingerprint, id).is_some() { 190 return Err(adapter_error(MycRelayAdapterErrorKind::Configuration)); 191 } 192 } 193 RelayProfileKind::Simulator => { 194 local.push(endpoint); 195 local_targets.push(target); 196 local_required |= required; 197 if local_ids.insert(fingerprint, id).is_some() { 198 return Err(adapter_error(MycRelayAdapterErrorKind::Configuration)); 199 } 200 } 201 RelayProfileKind::Device => { 202 return Err(adapter_error(MycRelayAdapterErrorKind::Configuration)); 203 } 204 _ => return Err(adapter_error(MycRelayAdapterErrorKind::Configuration)), 205 } 206 } 207 let mut groups = Vec::with_capacity(2); 208 add_ingress_group( 209 &mut groups, 210 RelayProfileKind::Public, 211 public, 212 public_targets, 213 public_ids, 214 public_required, 215 connect_timeout, 216 request_timeout, 217 "myc-runtime-public", 218 )?; 219 add_ingress_group( 220 &mut groups, 221 RelayProfileKind::Simulator, 222 local, 223 local_targets, 224 local_ids, 225 local_required, 226 connect_timeout, 227 request_timeout, 228 "myc-runtime-local", 229 )?; 230 if groups.is_empty() || groups.len() > 2 || !groups.iter().any(|group| group.required) { 231 return Err(adapter_error(MycRelayAdapterErrorKind::Configuration)); 232 } 233 Ok(Self { 234 groups: groups.into_boxed_slice(), 235 }) 236 } 237 238 pub(crate) fn group_count(&self) -> usize { 239 self.groups.len() 240 } 241 242 pub(crate) fn group_is_required(&self, index: usize) -> Option<bool> { 243 self.groups.get(index).map(|group| group.required) 244 } 245 246 pub(crate) async fn subscribe( 247 &self, 248 index: usize, 249 event_limit: u16, 250 deadline_unix_ms: u64, 251 checkpoints: &[SubscriptionCheckpoint], 252 ) -> Result<BoxSubscription, MycRelayAdapterError> { 253 let group = self 254 .groups 255 .get(index) 256 .ok_or_else(|| adapter_error(MycRelayAdapterErrorKind::Configuration))?; 257 let selector = FetchSelector::all() 258 .with_kinds(vec![NIP46_RPC_KIND]) 259 .map_err(|_| adapter_error(MycRelayAdapterErrorKind::Configuration))?; 260 let request = SubscriptionRequest::new( 261 group.request_id, 262 group.targets.clone(), 263 SubscriptionBounds::new(event_limit, deadline_unix_ms) 264 .map_err(|_| adapter_error(MycRelayAdapterErrorKind::Configuration))?, 265 ) 266 .map_err(|_| adapter_error(MycRelayAdapterErrorKind::Configuration))? 267 .with_selector(selector) 268 .with_checkpoints(checkpoints.iter().cloned()) 269 .map_err(|_| adapter_error(MycRelayAdapterErrorKind::Configuration))?; 270 group 271 .transport 272 .subscribe(request) 273 .await 274 .map_err(|_| adapter_error(MycRelayAdapterErrorKind::Execution)) 275 } 276 277 pub(crate) fn relay_id( 278 &self, 279 index: usize, 280 target: &TargetFingerprint, 281 ) -> Result<MycRateRelayId, MycRelayAdapterError> { 282 self.groups 283 .get(index) 284 .and_then(|group| group.relay_ids.get(target)) 285 .cloned() 286 .ok_or_else(|| adapter_error(MycRelayAdapterErrorKind::Target)) 287 } 288 } 289 290 #[allow(clippy::too_many_arguments)] 291 fn add_ingress_group( 292 groups: &mut Vec<IngressGroup>, 293 kind: RelayProfileKind, 294 endpoints: Vec<RelayEndpoint>, 295 targets: Vec<Target>, 296 relay_ids: BTreeMap<TargetFingerprint, MycRateRelayId>, 297 required: bool, 298 connect_timeout: u64, 299 request_timeout: u64, 300 request_id: &'static str, 301 ) -> Result<(), MycRelayAdapterError> { 302 if targets.is_empty() { 303 return Ok(()); 304 } 305 let transport = build_transport(kind, endpoints, connect_timeout, request_timeout)? 306 .ok_or_else(|| adapter_error(MycRelayAdapterErrorKind::Configuration))?; 307 let targets = 308 TargetSet::new(targets).map_err(|_| adapter_error(MycRelayAdapterErrorKind::Target))?; 309 if targets.len() != relay_ids.len() { 310 return Err(adapter_error(MycRelayAdapterErrorKind::Configuration)); 311 } 312 groups.push(IngressGroup { 313 transport, 314 targets, 315 relay_ids, 316 required, 317 request_id, 318 }); 319 Ok(()) 320 } 321 322 impl MycNostrDeliveryAdapter { 323 pub(crate) fn from_configuration( 324 configuration: &MycConfigDocumentV1, 325 ) -> Result<Self, MycRelayAdapterError> { 326 let relays = configuration 327 .normalized() 328 .pointer("/relays") 329 .and_then(serde_json::Value::as_array) 330 .ok_or_else(|| adapter_error(MycRelayAdapterErrorKind::Configuration))?; 331 let connect_timeout = 332 configuration_integer(configuration, "/transport/connect_deadline_ms")?; 333 let request_timeout = configuration_integer( 334 configuration, 335 "/transport/publish_retry/attempt_deadline_ms", 336 )?; 337 let mut public = Vec::new(); 338 let mut local = Vec::new(); 339 let mut definitions = Vec::with_capacity(relays.len()); 340 for relay in relays { 341 let id = relay 342 .pointer("/id") 343 .and_then(serde_json::Value::as_str) 344 .and_then(|value| MycDeliveryRelayId::new(value).ok()) 345 .ok_or_else(|| adapter_error(MycRelayAdapterErrorKind::Configuration))?; 346 let url = relay 347 .pointer("/url") 348 .and_then(serde_json::Value::as_str) 349 .ok_or_else(|| adapter_error(MycRelayAdapterErrorKind::Configuration))?; 350 let read = relay 351 .pointer("/read") 352 .and_then(serde_json::Value::as_bool) 353 .ok_or_else(|| adapter_error(MycRelayAdapterErrorKind::Configuration))?; 354 let write = relay 355 .pointer("/write") 356 .and_then(serde_json::Value::as_bool) 357 .ok_or_else(|| adapter_error(MycRelayAdapterErrorKind::Configuration))?; 358 let access = if write { 359 RelayAccess::ReadWrite 360 } else if read { 361 RelayAccess::ReadOnly 362 } else { 363 return Err(adapter_error(MycRelayAdapterErrorKind::Configuration)); 364 }; 365 let (kind, policy) = if url.starts_with("wss://") { 366 (RelayProfileKind::Public, RelayUrlPolicy::Public) 367 } else if configuration.profile() == MycConfigProfile::RepoLocal 368 && url.starts_with("ws://") 369 { 370 (RelayProfileKind::Simulator, RelayUrlPolicy::Local) 371 } else { 372 return Err(adapter_error(MycRelayAdapterErrorKind::Configuration)); 373 }; 374 let endpoint = RelayEndpoint::new(url, policy, access) 375 .map_err(|_| adapter_error(MycRelayAdapterErrorKind::Configuration))?; 376 match kind { 377 RelayProfileKind::Public => public.push(endpoint), 378 RelayProfileKind::Simulator => local.push(endpoint), 379 RelayProfileKind::Device => { 380 return Err(adapter_error(MycRelayAdapterErrorKind::Configuration)); 381 } 382 _ => return Err(adapter_error(MycRelayAdapterErrorKind::Configuration)), 383 } 384 definitions.push((id, url.to_owned(), kind)); 385 } 386 let public_transport = build_transport( 387 RelayProfileKind::Public, 388 public, 389 connect_timeout, 390 request_timeout, 391 )?; 392 let local_transport = build_transport( 393 RelayProfileKind::Simulator, 394 local, 395 connect_timeout, 396 request_timeout, 397 )?; 398 let mut targets = BTreeMap::new(); 399 for (id, url, kind) in definitions { 400 let transport = match kind { 401 RelayProfileKind::Public => public_transport.clone(), 402 RelayProfileKind::Simulator => local_transport.clone(), 403 RelayProfileKind::Device => None, 404 _ => None, 405 } 406 .ok_or_else(|| adapter_error(MycRelayAdapterErrorKind::Configuration))?; 407 let target = Target::nostr_relay(url) 408 .map_err(|_| adapter_error(MycRelayAdapterErrorKind::Target))?; 409 if targets 410 .insert(id, RelayTarget { transport, target }) 411 .is_some() 412 { 413 return Err(adapter_error(MycRelayAdapterErrorKind::Configuration)); 414 } 415 } 416 if targets.is_empty() { 417 return Err(adapter_error(MycRelayAdapterErrorKind::Configuration)); 418 } 419 Ok(Self { targets }) 420 } 421 422 pub(crate) async fn probe_required_relays( 423 configuration: &MycConfigDocumentV1, 424 deadline_unix_ms: u64, 425 ) -> Result<(), MycRelayAdapterError> { 426 let relays = configuration 427 .normalized() 428 .pointer("/relays") 429 .and_then(serde_json::Value::as_array) 430 .ok_or_else(|| adapter_error(MycRelayAdapterErrorKind::Configuration))?; 431 let connect_timeout = 432 configuration_integer(configuration, "/transport/connect_deadline_ms")?; 433 let request_timeout = configuration_integer( 434 configuration, 435 "/transport/publish_retry/attempt_deadline_ms", 436 )?; 437 let mut public = Vec::new(); 438 let mut public_targets = Vec::new(); 439 let mut local = Vec::new(); 440 let mut local_targets = Vec::new(); 441 for relay in relays.iter().filter(|relay| { 442 relay 443 .pointer("/required") 444 .and_then(serde_json::Value::as_bool) 445 == Some(true) 446 }) { 447 let url = relay 448 .pointer("/url") 449 .and_then(serde_json::Value::as_str) 450 .ok_or_else(|| adapter_error(MycRelayAdapterErrorKind::Configuration))?; 451 let target = Target::nostr_relay(url) 452 .map_err(|_| adapter_error(MycRelayAdapterErrorKind::Target))?; 453 let (kind, policy) = if url.starts_with("wss://") { 454 (RelayProfileKind::Public, RelayUrlPolicy::Public) 455 } else if configuration.profile() == MycConfigProfile::RepoLocal 456 && url.starts_with("ws://") 457 { 458 (RelayProfileKind::Simulator, RelayUrlPolicy::Local) 459 } else { 460 return Err(adapter_error(MycRelayAdapterErrorKind::Configuration)); 461 }; 462 let endpoint = RelayEndpoint::new(url, policy, RelayAccess::ReadOnly) 463 .map_err(|_| adapter_error(MycRelayAdapterErrorKind::Configuration))?; 464 match kind { 465 RelayProfileKind::Public => { 466 public.push(endpoint); 467 public_targets.push(target); 468 } 469 RelayProfileKind::Simulator => { 470 local.push(endpoint); 471 local_targets.push(target); 472 } 473 RelayProfileKind::Device => { 474 return Err(adapter_error(MycRelayAdapterErrorKind::Configuration)); 475 } 476 _ => return Err(adapter_error(MycRelayAdapterErrorKind::Configuration)), 477 } 478 } 479 if public_targets.is_empty() && local_targets.is_empty() { 480 return Err(adapter_error(MycRelayAdapterErrorKind::Configuration)); 481 } 482 probe_group( 483 RelayProfileKind::Public, 484 public, 485 public_targets, 486 connect_timeout, 487 request_timeout, 488 deadline_unix_ms, 489 "myc-doctor-public", 490 ) 491 .await?; 492 probe_group( 493 RelayProfileKind::Simulator, 494 local, 495 local_targets, 496 connect_timeout, 497 request_timeout, 498 deadline_unix_ms, 499 "myc-doctor-local", 500 ) 501 .await 502 } 503 } 504 505 async fn probe_group( 506 kind: RelayProfileKind, 507 endpoints: Vec<RelayEndpoint>, 508 targets: Vec<Target>, 509 connect_timeout: u64, 510 request_timeout: u64, 511 deadline_unix_ms: u64, 512 request_id: &'static str, 513 ) -> Result<(), MycRelayAdapterError> { 514 if targets.is_empty() { 515 return Ok(()); 516 } 517 let transport = build_transport(kind, endpoints, connect_timeout, request_timeout)? 518 .ok_or_else(|| adapter_error(MycRelayAdapterErrorKind::Configuration))?; 519 let target_set = 520 TargetSet::new(targets).map_err(|_| adapter_error(MycRelayAdapterErrorKind::Target))?; 521 let selector = FetchSelector::all() 522 .with_since_unix_seconds(u64::MAX) 523 .map_err(|_| adapter_error(MycRelayAdapterErrorKind::Configuration))?; 524 let request = FetchRequest::new( 525 request_id, 526 target_set, 527 FetchBounds::new(1, deadline_unix_ms) 528 .map_err(|_| adapter_error(MycRelayAdapterErrorKind::Configuration))?, 529 ) 530 .map_err(|_| adapter_error(MycRelayAdapterErrorKind::Configuration))? 531 .with_selector(selector); 532 let page = transport 533 .fetch(request) 534 .await 535 .map_err(|_| adapter_error(MycRelayAdapterErrorKind::Execution))?; 536 if page.target_outcomes().is_empty() 537 || page.target_outcomes().iter().any(|outcome| { 538 !matches!( 539 outcome.state(), 540 FetchTargetState::Complete | FetchTargetState::Partial 541 ) 542 }) 543 { 544 return Err(adapter_error(MycRelayAdapterErrorKind::Execution)); 545 } 546 Ok(()) 547 } 548 549 impl MycRelayAdapter for MycNostrDeliveryAdapter { 550 type Prepared = MycPreparedRelayDelivery; 551 552 fn prepare( 553 &self, 554 relay_id: &MycDeliveryRelayId, 555 request_id: String, 556 exact_event_bytes: &[u8], 557 deadline_unix_ms: u64, 558 ) -> Result<Self::Prepared, MycRelayAdapterError> { 559 let binding = self 560 .targets 561 .get(relay_id) 562 .ok_or_else(|| adapter_error(MycRelayAdapterErrorKind::Target))?; 563 let raw = core::str::from_utf8(exact_event_bytes) 564 .map_err(|_| adapter_error(MycRelayAdapterErrorKind::Payload))?; 565 let signed = Codec::decode_signed_event(raw) 566 .map_err(|_| adapter_error(MycRelayAdapterErrorKind::Payload))?; 567 if signed.raw_json().as_bytes() != exact_event_bytes { 568 return Err(adapter_error(MycRelayAdapterErrorKind::Payload)); 569 } 570 let targets = TargetSet::new(vec![binding.target.clone()]) 571 .map_err(|_| adapter_error(MycRelayAdapterErrorKind::Target))?; 572 let request = DeliveryRequest::new( 573 request_id, 574 DeliveryPayload::new(signed), 575 targets, 576 SatisfactionPolicy::new(SatisfactionClass::Accepted, TargetPolicy::all()), 577 deadline_unix_ms, 578 ) 579 .map_err(|_| adapter_error(MycRelayAdapterErrorKind::Preparation))?; 580 let prepared = binding 581 .transport 582 .prepare_delivery(request) 583 .map_err(|_| adapter_error(MycRelayAdapterErrorKind::Preparation))?; 584 Ok(MycPreparedRelayDelivery { 585 transport: binding.transport.clone(), 586 prepared, 587 }) 588 } 589 590 fn execute<'a>(&'a self, prepared: Self::Prepared) -> RelayExecutionFuture<'a> { 591 Box::pin(async move { 592 let MycPreparedRelayDelivery { 593 transport, 594 prepared, 595 } = prepared; 596 match transport.execute_prepared_delivery(prepared).await { 597 Err(_) => Ok(MycRelayExecutionOutcome::UnknownAcknowledgement), 598 Ok(receipt) => classify_receipt(receipt.target_receipts()), 599 } 600 }) 601 } 602 } 603 604 fn classify_receipt( 605 receipts: &[DeliveryTargetReceipt], 606 ) -> Result<MycRelayExecutionOutcome, MycRelayAdapterError> { 607 let [receipt] = receipts else { 608 return Err(adapter_error(MycRelayAdapterErrorKind::Execution)); 609 }; 610 Ok(match receipt.outcome().kind() { 611 DeliveryOutcomeKind::Accepted | DeliveryOutcomeKind::Delivered => { 612 MycRelayExecutionOutcome::Accepted 613 } 614 DeliveryOutcomeKind::Rejected => MycRelayExecutionOutcome::Rejected, 615 DeliveryOutcomeKind::Unavailable | DeliveryOutcomeKind::Failed => { 616 MycRelayExecutionOutcome::TransportFailed 617 } 618 }) 619 } 620 621 fn build_transport( 622 kind: RelayProfileKind, 623 endpoints: Vec<RelayEndpoint>, 624 connect_timeout: u64, 625 request_timeout: u64, 626 ) -> Result<Option<NostrTransport>, MycRelayAdapterError> { 627 if endpoints.is_empty() { 628 return Ok(None); 629 } 630 let maximum_connections = endpoints.len().min(8); 631 let profile = RelayProfile::explicit(kind, endpoints) 632 .map_err(|_| adapter_error(MycRelayAdapterErrorKind::Configuration))?; 633 let config = Config::from_profile(profile) 634 .with_timeouts(connect_timeout, request_timeout, connect_timeout) 635 .and_then(|config| config.with_max_connections(maximum_connections)) 636 .map_err(|_| adapter_error(MycRelayAdapterErrorKind::Configuration))?; 637 Ok(Some(NostrTransport::new(config))) 638 } 639 640 fn configuration_integer( 641 configuration: &MycConfigDocumentV1, 642 pointer: &str, 643 ) -> Result<u64, MycRelayAdapterError> { 644 configuration 645 .normalized() 646 .pointer(pointer) 647 .and_then(serde_json::Value::as_u64) 648 .ok_or_else(|| adapter_error(MycRelayAdapterErrorKind::Configuration)) 649 } 650 651 #[cfg(test)] 652 mod tests { 653 use std::error::Error as _; 654 655 use nostr::{EventBuilder, JsonUtil as _, Keys}; 656 use radroots_transport::{ 657 outcome::{DeliveryOutcome, Retryability}, 658 sink::DeliveryTargetReceipt, 659 }; 660 661 use super::*; 662 use crate::{MycConfigProfile, parse_myc_config_v1}; 663 664 const CONFIG: &str = include_str!("../contracts/services_hardening/config.v1.example.toml"); 665 666 fn receipt(outcome: DeliveryOutcome) -> DeliveryTargetReceipt { 667 DeliveryTargetReceipt::attempted( 668 Target::nostr_relay("wss://relay.example.test/").expect("target"), 669 outcome, 670 ) 671 } 672 673 #[test] 674 fn receipt_classification_preserves_all_four_durable_outcomes() { 675 for (outcome, expected) in [ 676 ( 677 DeliveryOutcome::accepted(), 678 MycRelayExecutionOutcome::Accepted, 679 ), 680 ( 681 DeliveryOutcome::delivered(), 682 MycRelayExecutionOutcome::Accepted, 683 ), 684 ( 685 DeliveryOutcome::rejected(), 686 MycRelayExecutionOutcome::Rejected, 687 ), 688 ( 689 DeliveryOutcome::unavailable(), 690 MycRelayExecutionOutcome::TransportFailed, 691 ), 692 ( 693 DeliveryOutcome::failed(Retryability::Retryable).expect("failed"), 694 MycRelayExecutionOutcome::TransportFailed, 695 ), 696 ] { 697 assert_eq!(classify_receipt(&[receipt(outcome)]).unwrap(), expected); 698 } 699 assert!(classify_receipt(&[]).is_err()); 700 assert!( 701 classify_receipt(&[ 702 receipt(DeliveryOutcome::accepted()), 703 receipt(DeliveryOutcome::accepted()), 704 ]) 705 .is_err() 706 ); 707 } 708 709 #[test] 710 fn exact_config_builds_without_network_io_and_errors_are_redacted() { 711 let configuration = parse_myc_config_v1(CONFIG.as_bytes(), MycConfigProfile::RepoLocal) 712 .expect("configuration"); 713 let adapter = MycNostrDeliveryAdapter::from_configuration(&configuration) 714 .expect("offline adapter construction"); 715 assert_eq!(adapter.targets.len(), 2); 716 let signing_keys = Keys::parse(&format!("02{}", "00".repeat(31))).expect("test keys"); 717 let event = EventBuilder::text_note("exact committed payload") 718 .sign_with_keys(&signing_keys) 719 .expect("signed event"); 720 let exact = event.as_json(); 721 let relay = MycDeliveryRelayId::new("primary").expect("relay"); 722 let prepared = adapter 723 .prepare(&relay, "offline-prepare".to_owned(), exact.as_bytes(), 1) 724 .expect("offline prepared delivery"); 725 assert_eq!( 726 format!("{prepared:?}"), 727 "MycPreparedRelayDelivery([redacted])" 728 ); 729 for kind in [ 730 MycRelayAdapterErrorKind::Configuration, 731 MycRelayAdapterErrorKind::Target, 732 MycRelayAdapterErrorKind::Payload, 733 MycRelayAdapterErrorKind::Preparation, 734 MycRelayAdapterErrorKind::Execution, 735 ] { 736 let error = adapter_error(kind); 737 assert_eq!(error.kind(), kind); 738 assert!(error.source().is_none()); 739 assert!(!format!("{error} {error:?}").contains("relay.example")); 740 } 741 } 742 }