lib

Core libraries for Radroots
git clone https://radroots.dev/git/lib.git
Log | Files | Refs | README

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 }