client.rs (37382B)
1 //! Client construction and lifecycle. 2 3 use std::{ 4 collections::{BTreeMap, BTreeSet}, 5 sync::{ 6 Arc, 7 atomic::{AtomicU8, Ordering}, 8 }, 9 }; 10 11 use radroots_signing::Signer; 12 use radroots_storage::Storage; 13 #[cfg(feature = "memory")] 14 use radroots_storage::{event::SourceGeneration, memory::MemoryStorage}; 15 use radroots_transport::{EventSink, EventSource}; 16 17 use crate::{ 18 Error, Result, 19 capability::{Availability, CapabilityId, CapabilityReport}, 20 }; 21 22 /// Cloneable handle to a composed Radroots client. 23 #[derive(Clone)] 24 pub struct Client { 25 inner: Arc<ClientInner>, 26 } 27 28 /// Explicit composition boundary for a [`Client`]. 29 #[derive(Default)] 30 pub struct ClientBuilder { 31 storage: Option<Arc<dyn Storage>>, 32 #[cfg(feature = "sync")] 33 sync_storage: Option<Arc<dyn radroots_sync::policy::SyncStorage>>, 34 signer: Option<Arc<dyn Signer>>, 35 source: Option<Arc<dyn EventSource>>, 36 sink: Option<Arc<dyn EventSink>>, 37 #[cfg(feature = "sync")] 38 sync: Option<radroots_sync::Engine>, 39 #[cfg(feature = "sync")] 40 host_sync: Option<crate::sync::HostPolicy>, 41 #[cfg(feature = "nostr")] 42 nostr: Option<crate::transport::NostrSlot>, 43 #[cfg(feature = "blossom")] 44 blossom: Option<crate::transport::BlossomSlot>, 45 capability_availability: BTreeMap<CapabilityId, Availability>, 46 explicitly_configured_capabilities: BTreeSet<CapabilityId>, 47 } 48 49 struct ClientInner { 50 storage: Arc<dyn Storage>, 51 signer: Option<Arc<dyn Signer>>, 52 source: Option<Arc<dyn EventSource>>, 53 sink: Option<Arc<dyn EventSink>>, 54 #[cfg(feature = "sync")] 55 sync: Option<radroots_sync::Engine>, 56 #[cfg(feature = "nostr")] 57 nostr: Option<crate::transport::NostrSlot>, 58 #[cfg(feature = "blossom")] 59 blossom: Option<crate::transport::BlossomSlot>, 60 capability_availability: BTreeMap<CapabilityId, Availability>, 61 explicitly_configured_capabilities: BTreeSet<CapabilityId>, 62 lifecycle: AtomicU8, 63 } 64 65 const OPEN: u8 = 0; 66 const CLOSING: u8 = 1; 67 const CLOSE_RETRY_REQUIRED: u8 = 2; 68 const CLOSED: u8 = 3; 69 70 impl ClientBuilder { 71 /// Creates an empty builder with no hidden storage, network, signing, or 72 /// runtime side effects. 73 #[must_use] 74 pub fn new() -> Self { 75 Self::default() 76 } 77 78 /// Creates a builder backed by deterministic in-process memory storage. 79 #[cfg(feature = "memory")] 80 #[must_use] 81 pub fn memory(generation: SourceGeneration) -> Self { 82 let storage = Arc::new(MemoryStorage::new(generation)); 83 let builder = Self::new().storage(storage.clone()); 84 #[cfg(feature = "sync")] 85 { 86 let mut builder = builder; 87 builder.sync_storage = Some(storage); 88 builder 89 } 90 #[cfg(not(feature = "sync"))] 91 builder 92 } 93 94 /// Creates the ordinary deterministic in-process memory configuration. 95 /// 96 /// This performs no I/O and is intended for ephemeral local clients whose 97 /// source-generation identity does not need to survive the process. Hosts 98 /// that persist cursors should call [`Self::memory`] with their own 99 /// generation instead. 100 #[cfg(feature = "memory")] 101 #[must_use] 102 pub fn memory_default() -> Self { 103 let storage = Arc::new(MemoryStorage::default()); 104 let builder = Self::new().storage(storage.clone()); 105 #[cfg(feature = "sync")] 106 { 107 let mut builder = builder; 108 builder.sync_storage = Some(storage); 109 builder 110 } 111 #[cfg(not(feature = "sync"))] 112 builder 113 } 114 115 /// Explicitly opens canonical SQLite storage from validated host-owned 116 /// configuration and returns a builder containing only the storage SPI. 117 #[cfg(feature = "sqlite")] 118 pub async fn sqlite(options: crate::storage::SqliteOptions) -> Result<Self> { 119 let storage = radroots_storage_sqlite::SqliteStorage::open(options) 120 .await 121 .map_err(Error::storage_open_failed)?; 122 let storage = Arc::new(storage); 123 let builder = Self::new() 124 .storage(storage.clone()) 125 .capability_availability(CapabilityId::PERSISTENT_STORAGE, Availability::Available); 126 #[cfg(feature = "sync")] 127 { 128 let mut builder = builder; 129 builder.sync_storage = Some(storage); 130 Ok(builder) 131 } 132 #[cfg(not(feature = "sync"))] 133 Ok(builder) 134 } 135 136 /// Injects the canonical storage capability. 137 #[must_use] 138 pub fn storage(mut self, storage: Arc<dyn Storage>) -> Self { 139 self.storage = Some(storage); 140 self 141 } 142 143 /// Injects an optional canonical signer capability. 144 #[must_use] 145 pub fn signer(mut self, signer: Arc<dyn Signer>) -> Self { 146 self.signer = Some(signer); 147 self 148 } 149 150 /// Installs a module-scoped signer provider over the canonical SPI and 151 /// records its presentation-independent runtime capability. 152 #[must_use] 153 pub fn signing(mut self, provider: crate::signing::Provider) -> Self { 154 let capability = match provider.mode() { 155 crate::signing::Mode::Local => Some(CapabilityId::LOCAL_SIGNING), 156 crate::signing::Mode::Nip46 => Some(CapabilityId::NIP46_SIGNING), 157 crate::signing::Mode::Host => None, 158 }; 159 self.signer = Some(provider.into_signer()); 160 if let Some(capability) = capability { 161 self.explicitly_configured_capabilities.insert(capability); 162 self.capability_availability 163 .insert(capability, Availability::Available); 164 } 165 self 166 } 167 168 /// Injects an optional inbound event source. 169 #[must_use] 170 pub fn source(mut self, source: Arc<dyn EventSource>) -> Self { 171 self.source = Some(source); 172 self 173 } 174 175 /// Injects an optional outbound event sink. 176 #[must_use] 177 pub fn sink(mut self, sink: Arc<dyn EventSink>) -> Self { 178 self.sink = Some(sink); 179 self 180 } 181 182 /// Installs one host-reconfigurable Nostr source and sink. 183 #[cfg(feature = "nostr")] 184 #[must_use] 185 pub fn nostr(mut self, slot: crate::transport::NostrSlot) -> Self { 186 self.source = Some(Arc::new(slot.clone())); 187 self.sink = Some(Arc::new(slot.clone())); 188 self.nostr = Some(slot); 189 self 190 } 191 192 /// Installs one host-reconfigurable Blossom HTTP adapter slot. 193 #[cfg(feature = "blossom")] 194 #[must_use] 195 pub fn blossom(mut self, slot: crate::transport::BlossomSlot) -> Self { 196 self.blossom = Some(slot); 197 self 198 } 199 200 /// Injects an explicitly composed synchronization engine. 201 #[cfg(feature = "sync")] 202 #[must_use] 203 pub fn sync_engine(mut self, sync: radroots_sync::Engine) -> Self { 204 self.sync = Some(sync); 205 self 206 } 207 208 /// Requests one SDK-composed engine using explicit native host policy. 209 /// 210 /// This is available only for SDK-created memory or SQLite storage, whose 211 /// complete synchronization capability is known without downcasting. 212 #[cfg(feature = "sync")] 213 #[must_use] 214 pub fn host_sync(mut self, policy: crate::sync::HostPolicy) -> Self { 215 self.host_sync = Some(policy); 216 self 217 } 218 219 /// Marks a specialized capability as configured and records its 220 /// host-observed initial availability without probing resources. 221 /// 222 /// Reports ignore this observation for capabilities that are not compiled 223 /// or configured. Runtime IDs are independent from Cargo feature names. 224 #[must_use] 225 pub fn capability_availability(mut self, id: CapabilityId, availability: Availability) -> Self { 226 self.explicitly_configured_capabilities.insert(id); 227 self.capability_availability.insert(id, availability); 228 self 229 } 230 231 /// Validates the selected capabilities and creates a client handle. 232 #[allow(unused_mut)] 233 pub fn build(mut self) -> Result<Client> { 234 let storage = self.storage.ok_or_else(Error::missing_storage)?; 235 #[cfg(feature = "sync")] 236 if let Some(policy) = self.host_sync { 237 let sync_storage = self 238 .sync_storage 239 .take() 240 .ok_or_else(Error::shared_operation_unavailable)?; 241 let (clock, ids, deadlines) = policy.composition(); 242 let mut builder = radroots_sync::Engine::builder(sync_storage, clock, ids, deadlines); 243 if let Some(source) = self.source.as_ref() { 244 builder = builder.source(Arc::clone(source)); 245 } 246 if let Some(sink) = self.sink.as_ref() { 247 builder = builder.sink(Arc::clone(sink)); 248 } 249 if let Some(signer) = self.signer.as_ref() { 250 builder = builder.signer(Arc::clone(signer)); 251 } 252 self.sync = Some(builder.build().map_err(Error::invalid_host_configuration)?); 253 } 254 Ok(Client { 255 inner: Arc::new(ClientInner { 256 storage, 257 signer: self.signer, 258 source: self.source, 259 sink: self.sink, 260 #[cfg(feature = "sync")] 261 sync: self.sync, 262 #[cfg(feature = "nostr")] 263 nostr: self.nostr, 264 #[cfg(feature = "blossom")] 265 blossom: self.blossom, 266 capability_availability: self.capability_availability, 267 explicitly_configured_capabilities: self.explicitly_configured_capabilities, 268 lifecycle: AtomicU8::new(OPEN), 269 }), 270 }) 271 } 272 } 273 274 impl Client { 275 /// Returns a deterministic capability report without probing resources or 276 /// performing filesystem, network, signing, or storage operations. 277 #[must_use] 278 #[allow(unused_mut)] 279 pub fn capabilities(&self) -> CapabilityReport { 280 let lifecycle_availability = match self.inner.lifecycle.load(Ordering::Acquire) { 281 OPEN => Availability::Available, 282 CLOSING | CLOSE_RETRY_REQUIRED => Availability::Degraded, 283 _ => Availability::Unavailable, 284 }; 285 let mut explicitly_configured = self.inner.explicitly_configured_capabilities.clone(); 286 let mut overrides = self.inner.capability_availability.clone(); 287 let mut source_configured = self.inner.source.is_some(); 288 let mut sink_configured = self.inner.sink.is_some(); 289 let mut source_availability = Availability::Unavailable; 290 let mut sink_availability = Availability::Unavailable; 291 #[cfg(feature = "nostr")] 292 if let Some(slot) = &self.inner.nostr { 293 match slot.relay_status() { 294 Some(status) => { 295 source_configured = !status.relays().is_empty(); 296 sink_configured = status 297 .relays() 298 .iter() 299 .any(|relay| relay.endpoint().access().can_write()); 300 source_availability = map_transport_availability(status.read_availability()); 301 sink_availability = map_transport_availability(status.write_availability()); 302 if source_configured { 303 explicitly_configured.insert(CapabilityId::NOSTR_FETCH); 304 overrides.insert(CapabilityId::NOSTR_FETCH, source_availability); 305 } 306 if sink_configured { 307 explicitly_configured.insert(CapabilityId::NOSTR_DELIVERY); 308 overrides.insert(CapabilityId::NOSTR_DELIVERY, sink_availability); 309 } 310 } 311 None => { 312 source_configured = false; 313 sink_configured = false; 314 } 315 } 316 } 317 crate::capability::report(crate::capability::Context { 318 storage: true, 319 signer: self.inner.signer.is_some(), 320 source: source_configured, 321 sink: sink_configured, 322 sync: self.sync_is_configured(), 323 source_availability, 324 sink_availability, 325 lifecycle_availability, 326 explicitly_configured: &explicitly_configured, 327 overrides: &overrides, 328 }) 329 } 330 331 /// Returns canonical backend status without exposing a backend handle. 332 pub async fn storage_status(&self) -> Result<crate::storage::Status> { 333 let storage = self.storage()?; 334 radroots_storage::BackupSource::status(storage) 335 .await 336 .map_err(Error::storage_inspection_failed) 337 } 338 339 /// Runs canonical integrity inspection without exposing backend internals. 340 pub async fn storage_integrity(&self) -> Result<crate::storage::IntegrityStatus> { 341 let storage = self.storage()?; 342 radroots_storage::BackupSource::integrity(storage) 343 .await 344 .map_err(Error::storage_inspection_failed) 345 } 346 347 /// Returns the injected canonical storage capability. 348 pub fn storage(&self) -> Result<&dyn Storage> { 349 self.require_open()?; 350 Ok(self.inner.storage.as_ref()) 351 } 352 353 /// Returns backend-neutral backup, restore, status, and integrity operations. 354 pub fn storage_operations(&self) -> Result<crate::storage::Operations<'_>> { 355 Ok(crate::storage::Operations::new(self.storage()?)) 356 } 357 358 /// Returns the injected signer, when outbound authoring is enabled. 359 pub fn signer(&self) -> Result<Option<&dyn Signer>> { 360 self.require_open()?; 361 Ok(self.inner.signer.as_deref()) 362 } 363 364 /// Returns focused high-level operations over the configured opaque signer. 365 pub fn signing(&self) -> Result<Option<crate::signing::Operations<'_>>> { 366 Ok(self.signer()?.map(crate::signing::Operations::new)) 367 } 368 369 /// Returns the injected inbound source, when pull is enabled. 370 pub fn source(&self) -> Result<Option<&dyn EventSource>> { 371 self.require_open()?; 372 Ok(self.inner.source.as_deref()) 373 } 374 375 /// Returns the injected outbound sink, when delivery is enabled. 376 pub fn sink(&self) -> Result<Option<&dyn EventSink>> { 377 self.require_open()?; 378 Ok(self.inner.sink.as_deref()) 379 } 380 381 /// Returns passive Nostr relay evidence when the concrete adapter is selected. 382 #[cfg(feature = "nostr")] 383 pub fn nostr_status(&self) -> Result<Option<crate::transport::RelayStatusReport>> { 384 self.require_open()?; 385 Ok(self 386 .inner 387 .nostr 388 .as_ref() 389 .and_then(crate::transport::NostrSlot::relay_status)) 390 } 391 392 /// Atomically replaces the active Nostr relay profile after validating it. 393 /// This does not perform DNS lookup or open a socket. 394 #[cfg(feature = "nostr")] 395 pub fn configure_nostr(&self, profile: crate::transport::RelayProfile) -> Result<()> { 396 self.require_open()?; 397 self.inner 398 .nostr 399 .as_ref() 400 .ok_or_else(Error::shared_operation_unavailable)? 401 .configure(profile) 402 } 403 404 /// Returns the explicitly composed Blossom adapter, when configured. 405 #[cfg(feature = "blossom")] 406 pub fn blossom(&self) -> Result<Option<&crate::transport::BlossomSlot>> { 407 self.require_open()?; 408 Ok(self.inner.blossom.as_ref()) 409 } 410 411 /// Atomically replaces the active Blossom profile without network I/O. 412 #[cfg(feature = "blossom")] 413 pub fn configure_blossom(&self, config: crate::transport::BlossomConfig) -> Result<()> { 414 self.require_open()?; 415 self.inner 416 .blossom 417 .as_ref() 418 .ok_or_else(Error::shared_operation_unavailable)? 419 .configure(config) 420 .map_err(Error::invalid_host_configuration) 421 } 422 423 /// Returns client-scoped canonical synchronization operations, when configured. 424 #[cfg(feature = "sync")] 425 pub fn sync(&self) -> Result<Option<crate::sync::Operations<'_>>> { 426 self.require_open()?; 427 Ok(self.inner.sync.as_ref().map(crate::sync::Operations::new)) 428 } 429 430 /// Returns farm commit operations when canonical synchronization is configured. 431 #[cfg(feature = "sync")] 432 pub fn farm(&self) -> Result<Option<crate::farm::Operations<'_>>> { 433 Ok(self.sync()?.map(crate::farm::Operations::new)) 434 } 435 436 /// Returns listing commit operations when canonical synchronization is configured. 437 #[cfg(feature = "sync")] 438 pub fn listing(&self) -> Result<Option<crate::listing::Operations<'_>>> { 439 Ok(self.sync()?.map(crate::listing::Operations::new)) 440 } 441 442 /// Returns trade operations when canonical synchronization is configured. 443 #[cfg(feature = "sync")] 444 pub fn trade(&self) -> Result<Option<crate::trade::Operations<'_>>> { 445 let storage = self.storage()?; 446 Ok(self 447 .sync()? 448 .map(|sync| crate::trade::Operations::new(storage, sync))) 449 } 450 451 /// Returns whether explicit close completed successfully or reached the 452 /// lower storage commit point. 453 #[must_use] 454 pub fn is_closed(&self) -> bool { 455 self.inner.lifecycle.load(Ordering::Acquire) == CLOSED 456 } 457 458 /// Explicitly closes active storage resources across every client clone. 459 /// 460 /// Dropping this future before its first poll has no effect. Cancellation 461 /// after close begins leaves the client unavailable and permits an 462 /// explicit retry; it never reports rollback. Repeated completed close is 463 /// idempotent. No worker, executor, or blocking `Drop` path is installed. 464 pub async fn close(&self) -> Result<()> { 465 loop { 466 match self.inner.lifecycle.load(Ordering::Acquire) { 467 CLOSED => return Ok(()), 468 CLOSING => return Err(Error::close_in_progress()), 469 state @ (OPEN | CLOSE_RETRY_REQUIRED) => { 470 if self 471 .inner 472 .lifecycle 473 .compare_exchange(state, CLOSING, Ordering::AcqRel, Ordering::Acquire) 474 .is_ok() 475 { 476 break; 477 } 478 } 479 _ => return Err(Error::client_closed()), 480 } 481 } 482 483 let attempt = CloseAttempt::new(Arc::clone(&self.inner)); 484 let close_result = radroots_storage::BackupSource::close(self.inner.storage.as_ref()).await; 485 attempt.complete(); 486 close_result 487 .map(|_| ()) 488 .map_err(Error::storage_close_failed) 489 } 490 491 fn require_open(&self) -> Result<()> { 492 match self.inner.lifecycle.load(Ordering::Acquire) { 493 OPEN => Ok(()), 494 CLOSING | CLOSE_RETRY_REQUIRED => Err(Error::client_closing()), 495 _ => Err(Error::client_closed()), 496 } 497 } 498 499 #[cfg(feature = "sync")] 500 fn sync_is_configured(&self) -> bool { 501 self.inner.sync.is_some() 502 } 503 504 #[cfg(not(feature = "sync"))] 505 fn sync_is_configured(&self) -> bool { 506 false 507 } 508 } 509 510 #[cfg(feature = "nostr")] 511 const fn map_transport_availability( 512 availability: radroots_transport::capability::Availability, 513 ) -> Availability { 514 match availability { 515 radroots_transport::capability::Availability::Available => Availability::Available, 516 radroots_transport::capability::Availability::Degraded => Availability::Degraded, 517 radroots_transport::capability::Availability::Unavailable => Availability::Unavailable, 518 } 519 } 520 521 struct CloseAttempt { 522 inner: Arc<ClientInner>, 523 completed: bool, 524 } 525 526 impl CloseAttempt { 527 fn new(inner: Arc<ClientInner>) -> Self { 528 Self { 529 inner, 530 completed: false, 531 } 532 } 533 534 fn complete(mut self) { 535 self.inner.lifecycle.store(CLOSED, Ordering::Release); 536 self.completed = true; 537 } 538 } 539 540 impl Drop for CloseAttempt { 541 fn drop(&mut self) { 542 if !self.completed { 543 self.inner 544 .lifecycle 545 .store(CLOSE_RETRY_REQUIRED, Ordering::Release); 546 } 547 } 548 } 549 550 impl std::fmt::Debug for Client { 551 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { 552 let mut debug = formatter.debug_struct("Client"); 553 debug 554 .field("signer", &self.inner.signer.is_some()) 555 .field("source", &self.inner.source.is_some()) 556 .field("sink", &self.inner.sink.is_some()); 557 #[cfg(feature = "blossom")] 558 debug.field("blossom", &self.inner.blossom.is_some()); 559 debug 560 .field("closed", &self.is_closed()) 561 .finish_non_exhaustive() 562 } 563 } 564 565 impl std::fmt::Debug for ClientBuilder { 566 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { 567 let mut debug = formatter.debug_struct("ClientBuilder"); 568 debug 569 .field("storage", &self.storage.is_some()) 570 .field("signer", &self.signer.is_some()) 571 .field("source", &self.source.is_some()) 572 .field("sink", &self.sink.is_some()); 573 #[cfg(feature = "blossom")] 574 debug.field("blossom", &self.blossom.is_some()); 575 debug.finish_non_exhaustive() 576 } 577 } 578 579 #[cfg(all(test, feature = "memory"))] 580 mod tests { 581 use super::*; 582 use radroots_signing::{ 583 Error as SigningError, SignReceipt, SignRequest, SignerStatus, error::Kind, 584 signer::BoxFuture as SigningFuture, 585 }; 586 use radroots_transport::{ 587 DeliveryReceipt, DeliveryRequest, Error as TransportError, FetchPage, FetchRequest, 588 SinkFailure, SinkStatus, SourceStatus, outcome::Retryability, 589 source::BoxFuture as TransportFuture, 590 }; 591 use std::{ 592 future::Future, 593 sync::Arc, 594 task::{Context, Poll, Wake, Waker}, 595 }; 596 597 struct TestSource; 598 struct TestSink; 599 struct TestSigner; 600 601 impl EventSource for TestSource { 602 fn status(&self) -> TransportFuture<'_, std::result::Result<SourceStatus, TransportError>> { 603 Box::pin(async { Err(TransportError::UnsupportedOperation) }) 604 } 605 606 fn fetch( 607 &self, 608 _request: FetchRequest, 609 ) -> TransportFuture<'_, std::result::Result<FetchPage, TransportError>> { 610 Box::pin(async { Err(TransportError::UnsupportedOperation) }) 611 } 612 } 613 614 impl EventSink for TestSink { 615 fn status(&self) -> TransportFuture<'_, std::result::Result<SinkStatus, TransportError>> { 616 Box::pin(async { Err(TransportError::UnsupportedOperation) }) 617 } 618 619 fn deliver( 620 &self, 621 request: DeliveryRequest, 622 ) -> TransportFuture<'_, std::result::Result<DeliveryReceipt, SinkFailure>> { 623 Box::pin(async move { 624 Err(SinkFailure::for_request( 625 &request, 626 "test_sink_unavailable", 627 Retryability::Terminal, 628 None, 629 None, 630 Vec::new(), 631 ) 632 .expect("test sink failure")) 633 }) 634 } 635 } 636 637 impl Signer for TestSigner { 638 fn status(&self) -> SigningFuture<'_, std::result::Result<SignerStatus, SigningError>> { 639 Box::pin(async { Err(SigningError::new(Kind::InternalError)) }) 640 } 641 642 fn sign( 643 &self, 644 _request: SignRequest, 645 ) -> SigningFuture<'_, std::result::Result<SignReceipt, SigningError>> { 646 Box::pin(async { Err(SigningError::new(Kind::InternalError)) }) 647 } 648 } 649 650 fn generation() -> SourceGeneration { 651 SourceGeneration::new([1; 32]).expect("non-zero generation") 652 } 653 654 #[test] 655 fn missing_storage_fails_and_signer_only_composition_is_valid() { 656 assert!(matches!( 657 ClientBuilder::new().build(), 658 Err(error) if error.kind() == crate::error::ErrorKind::MissingStorage 659 )); 660 let signer_only = ClientBuilder::memory(generation()) 661 .signer(Arc::new(TestSigner)) 662 .build() 663 .expect("HTTP-only signing does not require a relay sink"); 664 assert!(signer_only.signer().expect("signer").is_some()); 665 assert!(signer_only.signing().expect("signing operations").is_some()); 666 assert!(signer_only.sink().expect("sink").is_none()); 667 } 668 669 #[test] 670 fn memory_local_source_only_and_sink_only_compositions_are_explicit() { 671 let local = ClientBuilder::memory(generation()).build().expect("local"); 672 assert!(local.source().expect("source capability").is_none()); 673 assert!(local.sink().expect("sink capability").is_none()); 674 assert!(local.signer().expect("signer capability").is_none()); 675 676 let source = ClientBuilder::memory(generation()) 677 .source(Arc::new(TestSource)) 678 .build() 679 .expect("source-only"); 680 assert!(source.source().expect("source capability").is_some()); 681 assert!(source.sink().expect("sink capability").is_none()); 682 683 let sink = ClientBuilder::memory(generation()) 684 .sink(Arc::new(TestSink)) 685 .build() 686 .expect("sink-only"); 687 assert!(sink.source().expect("source capability").is_none()); 688 assert!(sink.sink().expect("sink capability").is_some()); 689 } 690 691 #[cfg(feature = "nostr")] 692 #[test] 693 fn nostr_capabilities_require_directional_profile_authority_and_evidence() { 694 let slot = crate::transport::NostrSlot::new(); 695 slot.configure( 696 crate::transport::RelayProfile::explicit( 697 crate::transport::RelayProfileKind::Public, 698 [crate::transport::RelayEndpoint::new( 699 "wss://radroots.org", 700 crate::transport::RelayUrlPolicy::Public, 701 crate::transport::RelayAccess::ReadOnly, 702 ) 703 .expect("read-only endpoint")], 704 ) 705 .expect("read-only public profile"), 706 ) 707 .expect("configure slot"); 708 let client = ClientBuilder::memory(generation()) 709 .nostr(slot) 710 .build() 711 .expect("client"); 712 713 let capabilities = client.capabilities(); 714 let fetch = capabilities 715 .get(CapabilityId::NOSTR_FETCH) 716 .expect("fetch capability"); 717 let delivery = capabilities 718 .get(CapabilityId::NOSTR_DELIVERY) 719 .expect("delivery capability"); 720 assert!(fetch.is_configured()); 721 assert_eq!(fetch.availability(), Availability::Unavailable); 722 assert!(!delivery.is_configured()); 723 assert_eq!(delivery.availability(), Availability::Unavailable); 724 725 client 726 .configure_nostr( 727 crate::transport::RelayProfile::explicit( 728 crate::transport::RelayProfileKind::Public, 729 [crate::transport::RelayEndpoint::new( 730 "wss://write.example", 731 crate::transport::RelayUrlPolicy::Public, 732 crate::transport::RelayAccess::ReadWrite, 733 ) 734 .expect("writable endpoint")], 735 ) 736 .expect("writable profile"), 737 ) 738 .expect("reconfigure"); 739 let capabilities = client.capabilities(); 740 let delivery = capabilities 741 .get(CapabilityId::NOSTR_DELIVERY) 742 .expect("delivery capability"); 743 assert!(delivery.is_configured()); 744 assert_eq!(delivery.availability(), Availability::Unavailable); 745 } 746 747 #[test] 748 fn signer_with_sink_is_valid_and_diagnostics_are_capability_only() { 749 let client = ClientBuilder::memory(generation()) 750 .sink(Arc::new(TestSink)) 751 .signer(Arc::new(TestSigner)) 752 .build() 753 .expect("outbound client"); 754 assert!(client.signer().expect("signer capability").is_some()); 755 let expected = if cfg!(feature = "blossom") { 756 "Client { signer: true, source: false, sink: true, blossom: false, closed: false, .. }" 757 } else { 758 "Client { signer: true, source: false, sink: true, closed: false, .. }" 759 }; 760 assert_eq!(format!("{client:?}"), expected); 761 } 762 763 #[cfg(all(feature = "sync", feature = "nip46"))] 764 #[test] 765 fn canonical_operation_accessors_and_builder_modes_are_complete() { 766 let builder = ClientBuilder::memory_default() 767 .source(Arc::new(TestSource)) 768 .host_sync(crate::sync::HostPolicy::default()); 769 assert!(format!("{builder:?}").contains("storage: true")); 770 let client = builder.build().expect("sync client"); 771 assert!(client.sync().expect("sync").is_some()); 772 assert!(client.farm().expect("farm").is_some()); 773 assert!(client.listing().expect("listing").is_some()); 774 assert!(client.trade().expect("trade").is_some()); 775 assert!( 776 format!( 777 "{:?}", 778 client.storage_operations().expect("storage operations") 779 ) 780 .contains("borrowed canonical storage") 781 ); 782 assert!( 783 format!("{:?}", client.sync().expect("sync").expect("operations")) 784 .contains("borrowed canonical engine") 785 ); 786 787 let host = ClientBuilder::memory_default() 788 .signing(crate::signing::Provider::host(Arc::new(TestSigner))) 789 .build() 790 .expect("host signer"); 791 assert!( 792 !host 793 .capabilities() 794 .get(CapabilityId::NIP46_SIGNING) 795 .expect("nip46") 796 .is_configured() 797 ); 798 799 let nip46 = ClientBuilder::memory_default() 800 .signing(crate::signing::Provider::nip46(Arc::new(TestSigner))) 801 .build() 802 .expect("nip46 signer"); 803 assert!( 804 nip46 805 .capabilities() 806 .get(CapabilityId::NIP46_SIGNING) 807 .expect("nip46") 808 .is_configured() 809 ); 810 } 811 812 #[cfg(feature = "local-signing")] 813 #[test] 814 fn module_scoped_signing_provider_configures_the_matching_capability() { 815 let signer = radroots_nostr::signing::LocalSigner::generate().expect("local signer"); 816 let client = ClientBuilder::memory(generation()) 817 .sink(Arc::new(TestSink)) 818 .signing(crate::signing::Provider::local(signer)) 819 .build() 820 .expect("client"); 821 let status = client 822 .capabilities() 823 .get(CapabilityId::LOCAL_SIGNING) 824 .expect("local signing"); 825 assert!(status.is_compiled()); 826 assert!(status.is_configured()); 827 assert_eq!(status.availability(), Availability::Available); 828 } 829 830 #[test] 831 fn close_is_clone_shared_idempotent_and_rejects_later_capability_access() { 832 let client = ClientBuilder::memory(generation()).build().expect("client"); 833 let clone = client.clone(); 834 assert!(!client.is_closed()); 835 block_on(client.close()).expect("first close"); 836 assert!(clone.is_closed()); 837 block_on(clone.close()).expect("repeated close"); 838 assert!(matches!( 839 clone.storage(), 840 Err(error) if error.kind() == crate::error::ErrorKind::ClientClosed 841 )); 842 assert!(matches!( 843 client.source(), 844 Err(error) if error.kind() == crate::error::ErrorKind::ClientClosed 845 )); 846 } 847 848 #[test] 849 fn close_cancellation_boundaries_are_explicit_and_retryable() { 850 let client = ClientBuilder::memory(generation()).build().expect("client"); 851 let unpolled = client.close(); 852 drop(unpolled); 853 assert!(client.storage().is_ok()); 854 855 client.inner.lifecycle.store(CLOSING, Ordering::Release); 856 let attempt = CloseAttempt::new(Arc::clone(&client.inner)); 857 drop(attempt); 858 assert!(matches!( 859 client.storage(), 860 Err(error) if error.kind() == crate::error::ErrorKind::ClientClosing 861 )); 862 block_on(client.close()).expect("retry close"); 863 assert!(client.is_closed()); 864 } 865 866 #[test] 867 fn concurrent_clones_converge_on_one_closed_state() { 868 let client = ClientBuilder::memory(generation()).build().expect("client"); 869 let first = client.clone(); 870 let second = client.clone(); 871 let outcomes = std::thread::scope(|scope| { 872 let first_close = scope.spawn(move || block_on(first.close())); 873 let second_close = scope.spawn(move || block_on(second.close())); 874 [ 875 first_close.join().expect("first thread"), 876 second_close.join().expect("second thread"), 877 ] 878 }); 879 assert!(outcomes.iter().all(|outcome| { 880 outcome.is_ok() 881 || matches!( 882 outcome, 883 Err(error) 884 if error.kind() == crate::error::ErrorKind::CloseInProgress 885 ) 886 })); 887 if !client.is_closed() { 888 block_on(client.close()).expect("finish close"); 889 } 890 assert!(client.is_closed()); 891 } 892 893 #[test] 894 fn capability_reports_separate_configuration_degradation_and_lifecycle() { 895 let local = ClientBuilder::memory(generation()) 896 .capability_availability(CapabilityId::CANONICAL_STORAGE, Availability::Degraded) 897 .build() 898 .expect("client"); 899 let report = local.capabilities(); 900 let storage = report 901 .get(CapabilityId::CANONICAL_STORAGE) 902 .expect("storage"); 903 assert!(storage.is_compiled()); 904 assert!(storage.is_configured()); 905 assert_eq!(storage.availability(), Availability::Degraded); 906 907 let signing = report.get(CapabilityId::LOCAL_SIGNING).expect("signing"); 908 assert!(!signing.is_configured()); 909 assert!(matches!( 910 signing.availability(), 911 Availability::Unavailable | Availability::Unsupported 912 )); 913 914 block_on(local.close()).expect("close"); 915 assert_eq!( 916 local 917 .capabilities() 918 .get(CapabilityId::BACKUP_RESTORE) 919 .expect("backup") 920 .availability(), 921 Availability::Unavailable 922 ); 923 } 924 925 #[test] 926 fn memory_storage_status_integrity_and_lifecycle_use_native_contracts() { 927 use radroots_storage::status::{ 928 IntegrityHealth, ShutdownState, StorageBackend, StorageOpenMode, WriterPolicy, 929 }; 930 931 let client = ClientBuilder::memory(generation()).build().expect("client"); 932 let status = block_on(client.storage_status()).expect("status"); 933 assert_eq!(status.backend(), StorageBackend::Memory); 934 assert_eq!(status.open_mode(), StorageOpenMode::Create); 935 assert_eq!(status.writer_policy(), WriterPolicy::NoWriter); 936 assert_eq!(status.shutdown(), ShutdownState::Open); 937 assert_eq!( 938 block_on(client.storage_integrity()) 939 .expect("integrity") 940 .health(), 941 IntegrityHealth::Healthy 942 ); 943 block_on(client.close()).expect("close"); 944 assert_eq!( 945 block_on(client.storage_status()) 946 .expect_err("closed") 947 .kind(), 948 crate::error::ErrorKind::ClientClosed 949 ); 950 } 951 952 #[cfg(feature = "sqlite")] 953 #[tokio::test] 954 async fn sqlite_builder_exposes_native_status_integrity_and_lifecycle() { 955 use radroots_storage::status::{ 956 IntegrityHealth, ShutdownState, StorageBackend, StorageOpenMode, WriterPolicy, 957 }; 958 959 let directory = tempfile::tempdir().expect("temporary directory"); 960 let paths = crate::storage::SqlitePaths::from_directory(directory.path()).expect("paths"); 961 let options = 962 crate::storage::SqliteOptions::new(paths, crate::storage::SqliteOpenMode::Create) 963 .with_source_generation(generation(), 1) 964 .expect("source generation"); 965 let client = ClientBuilder::sqlite(options) 966 .await 967 .expect("open builder") 968 .build() 969 .expect("client"); 970 let status = client.storage_status().await.expect("status"); 971 assert_eq!(status.backend(), StorageBackend::Sqlite); 972 assert_eq!(status.open_mode(), StorageOpenMode::Create); 973 assert_eq!(status.writer_policy(), WriterPolicy::AdvisoryProcessLock); 974 assert_eq!(status.shutdown(), ShutdownState::Open); 975 assert!(status.wal_enabled()); 976 assert_ne!(status.busy_timeout_ms(), 0); 977 assert_eq!( 978 client 979 .storage_integrity() 980 .await 981 .expect("integrity") 982 .health(), 983 IntegrityHealth::Unknown 984 ); 985 client.close().await.expect("close"); 986 assert!(client.is_closed()); 987 } 988 989 struct ThreadWaker; 990 991 impl Wake for ThreadWaker { 992 fn wake(self: Arc<Self>) { 993 std::thread::current().unpark(); 994 } 995 } 996 997 fn block_on<F: Future>(future: F) -> F::Output { 998 let waker = Waker::from(Arc::new(ThreadWaker)); 999 let mut context = Context::from_waker(&waker); 1000 let mut future = Box::pin(future); 1001 loop { 1002 match future.as_mut().poll(&mut context) { 1003 Poll::Ready(output) => return output, 1004 Poll::Pending => std::thread::park(), 1005 } 1006 } 1007 } 1008 }