app

Local-first trade for farms and co-ops
git clone https://radroots.dev/git/app.git
Log | Files | Refs | README | LICENSE

client.rs (11442B)


      1 use core::fmt;
      2 use std::sync::atomic::{AtomicU64, Ordering};
      3 use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH};
      4 
      5 use harvestcircle_application::{
      6     BoxFuture, MAX_CONFIGURED_RELAYS, NostrClient, ProfileFetchResult,
      7 };
      8 use harvestcircle_domain::{PublicKey, SafeError, SafeErrorCode, SafeMessage, select_latest_kind0};
      9 use radroots_transport::outcome::FetchTargetState;
     10 use radroots_transport::source::{FetchBounds, FetchSelector};
     11 use radroots_transport::{EventSource, FetchRequest, TargetSet};
     12 use radroots_transport_nostr::{Config, NostrTransport, RelayEndpoint, RelayProfile};
     13 
     14 const MAX_PROFILE_EVENTS_PER_FETCH: u16 = 64;
     15 
     16 pub struct SdkNostrClient {
     17     timeout: Duration,
     18     next_request: AtomicU64,
     19 }
     20 
     21 impl fmt::Debug for SdkNostrClient {
     22     fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
     23         formatter
     24             .debug_struct("SdkNostrClient")
     25             .field("timeout", &self.timeout)
     26             .field("transport", &"[sealed]")
     27             .finish()
     28     }
     29 }
     30 
     31 impl SdkNostrClient {
     32     #[must_use]
     33     pub const fn new(timeout: Duration) -> Self {
     34         Self {
     35             timeout,
     36             next_request: AtomicU64::new(1),
     37         }
     38     }
     39 
     40     fn request_id(&self) -> Result<String, SafeError> {
     41         let sequence = self
     42             .next_request
     43             .fetch_update(Ordering::AcqRel, Ordering::Acquire, |value| {
     44                 value.checked_add(1)
     45             })
     46             .map_err(|_| relay_connection_failed())?;
     47         Ok(format!("harvest-profile-{sequence}"))
     48     }
     49 }
     50 
     51 impl NostrClient for SdkNostrClient {
     52     fn fetch_profile<'a>(
     53         &'a self,
     54         public_key: PublicKey,
     55         relays: &'a [RelayEndpoint],
     56         deadline: Instant,
     57     ) -> BoxFuture<'a, Result<ProfileFetchResult, SafeError>> {
     58         Box::pin(async move {
     59             if relays.is_empty() || relays.len() > MAX_CONFIGURED_RELAYS {
     60                 return Err(invalid_relay_configuration());
     61             }
     62             let readable = relays
     63                 .iter()
     64                 .filter(|endpoint| endpoint.access().can_read())
     65                 .cloned()
     66                 .collect::<Vec<_>>();
     67             if readable.is_empty() {
     68                 return Err(invalid_relay_configuration());
     69             }
     70 
     71             let profile = RelayProfile::explicit(profile_kind(relays)?, relays.iter().cloned())
     72                 .map_err(|_| invalid_relay_configuration())?;
     73             let timeout_ms = u64::try_from(self.timeout.as_millis())
     74                 .ok()
     75                 .filter(|value| *value > 0)
     76                 .ok_or_else(invalid_relay_configuration)?;
     77             let config = Config::from_profile(profile)
     78                 .with_timeouts(timeout_ms, timeout_ms, timeout_ms)
     79                 .and_then(|config| config.with_max_connections(readable.len()))
     80                 .map_err(|_| invalid_relay_configuration())?;
     81             let transport = NostrTransport::new(config);
     82 
     83             let targets = readable
     84                 .iter()
     85                 .map(|endpoint| endpoint.url().to_target())
     86                 .collect::<Result<Vec<_>, _>>()
     87                 .map_err(|_| invalid_relay_configuration())?;
     88             let targets = TargetSet::new(targets).map_err(|_| invalid_relay_configuration())?;
     89             let remaining = deadline.saturating_duration_since(Instant::now());
     90             if remaining.is_zero() {
     91                 return Err(relay_connection_failed());
     92             }
     93             let deadline_unix_ms = current_unix_ms()?
     94                 .checked_add(
     95                     u64::try_from(remaining.as_millis()).map_err(|_| relay_connection_failed())?,
     96                 )
     97                 .ok_or_else(relay_connection_failed)?;
     98             let bounds = FetchBounds::new(MAX_PROFILE_EVENTS_PER_FETCH, deadline_unix_ms)
     99                 .map_err(|_| invalid_relay_configuration())?;
    100             let author = radroots_identity::PublicKey::from_bytes(*public_key.as_bytes())
    101                 .map_err(|_| invalid_relay_configuration())?;
    102             let selector = FetchSelector::all()
    103                 .with_kinds(vec![0])
    104                 .and_then(|selector| selector.with_authors(vec![author]))
    105                 .map_err(|_| invalid_relay_configuration())?;
    106             let request = FetchRequest::new(self.request_id()?, targets, bounds)
    107                 .map_err(|_| invalid_relay_configuration())?
    108                 .with_selector(selector);
    109             let page = transport
    110                 .fetch(request)
    111                 .await
    112                 .map_err(|_| relay_connection_failed())?;
    113 
    114             let successful = page
    115                 .target_outcomes()
    116                 .iter()
    117                 .filter(|outcome| outcome.state() == FetchTargetState::Complete)
    118                 .count();
    119             if successful == 0 {
    120                 return Err(relay_connection_failed());
    121             }
    122             let mut candidates = Vec::with_capacity(page.events().len());
    123             for observed in page.events() {
    124                 candidates.push(crate::parse_verified_kind0(
    125                     observed.event().raw_json(),
    126                     public_key,
    127                 )?);
    128             }
    129             let candidate = select_latest_kind0(candidates);
    130             if successful == readable.len() {
    131                 Ok(ProfileFetchResult::complete(candidate))
    132             } else {
    133                 Ok(ProfileFetchResult::partial(candidate))
    134             }
    135         })
    136     }
    137 }
    138 
    139 fn profile_kind(
    140     relays: &[RelayEndpoint],
    141 ) -> Result<radroots_transport_nostr::RelayProfileKind, SafeError> {
    142     let has_local = relays
    143         .iter()
    144         .any(|relay| relay.policy() == radroots_transport_nostr::RelayUrlPolicy::Local);
    145     let has_private = relays
    146         .iter()
    147         .any(|relay| relay.policy() == radroots_transport_nostr::RelayUrlPolicy::PrivateNetwork);
    148     let has_public = relays
    149         .iter()
    150         .any(|relay| relay.policy() == radroots_transport_nostr::RelayUrlPolicy::Public);
    151     match (has_local, has_private, has_public) {
    152         (true, false, false) => Ok(radroots_transport_nostr::RelayProfileKind::Simulator),
    153         (false, false, true) => Ok(radroots_transport_nostr::RelayProfileKind::Public),
    154         (false, true, _) => Ok(radroots_transport_nostr::RelayProfileKind::Device),
    155         _ => Err(invalid_relay_configuration()),
    156     }
    157 }
    158 
    159 fn current_unix_ms() -> Result<u64, SafeError> {
    160     SystemTime::now()
    161         .duration_since(UNIX_EPOCH)
    162         .ok()
    163         .and_then(|duration| u64::try_from(duration.as_millis()).ok())
    164         .ok_or_else(relay_connection_failed)
    165 }
    166 
    167 const fn invalid_relay_configuration() -> SafeError {
    168     SafeError::new(
    169         SafeErrorCode::InvalidRelayConfiguration,
    170         SafeMessage::new("The Nostr relay configuration is invalid."),
    171     )
    172 }
    173 
    174 const fn relay_connection_failed() -> SafeError {
    175     SafeError::new(
    176         SafeErrorCode::RelayConnectionFailed,
    177         SafeMessage::new("The Nostr relays could not be reached."),
    178     )
    179 }
    180 
    181 #[cfg(test)]
    182 mod tests {
    183     use std::time::Duration;
    184 
    185     use harvestcircle_application::{NostrClient, RelayFetchCompleteness};
    186     use harvestcircle_domain::{PublicKey, SafeErrorCode};
    187     use nostr::{EventBuilder, Keys, Metadata};
    188     use nostr_relay_builder::MockRelay;
    189     use nostr_sdk::Client;
    190     use radroots_transport_nostr::{RelayAccess, RelayEndpoint, RelayUrlPolicy};
    191 
    192     use crate::SdkNostrClient;
    193 
    194     #[tokio::test]
    195     async fn governed_transport_fetches_verified_profile_from_local_relay() {
    196         let relay = MockRelay::run().await.expect("local relay");
    197         let relay_url = relay.url().await;
    198         let keys = Keys::generate();
    199         let publisher = Client::new(keys.clone());
    200         publisher.add_relay(relay_url.clone()).await.expect("relay");
    201         publisher.connect().await;
    202         publisher.wait_for_connection(Duration::from_secs(2)).await;
    203         publisher
    204             .send_event_builder(EventBuilder::metadata(
    205                 &Metadata::new().name("Farmer").display_name("Farm Identity"),
    206             ))
    207             .await
    208             .expect("publish metadata");
    209 
    210         let adapter = SdkNostrClient::new(Duration::from_secs(2));
    211         let public_key = PublicKey::from_bytes(keys.public_key().to_bytes()).expect("public key");
    212         let fetched = adapter
    213             .fetch_profile(
    214                 public_key,
    215                 &[endpoint(relay_url.as_str(), RelayUrlPolicy::Local)],
    216                 std::time::Instant::now() + Duration::from_secs(2),
    217             )
    218             .await
    219             .expect("profile");
    220         let (profile, completeness) = fetched.into_parts();
    221         assert_eq!(
    222             profile
    223                 .expect("published profile")
    224                 .metadata()
    225                 .preferred_name(),
    226             Some("Farm Identity")
    227         );
    228         assert_eq!(completeness, RelayFetchCompleteness::Complete);
    229         publisher.shutdown().await;
    230         relay.shutdown();
    231     }
    232 
    233     #[tokio::test]
    234     async fn governed_transport_rejects_invalid_profiles_before_network_access() {
    235         let adapter = SdkNostrClient::new(Duration::from_millis(10));
    236         let public_key = PublicKey::from_bytes([7; 32]).expect("public key");
    237         let empty = adapter
    238             .fetch_profile(
    239                 public_key,
    240                 &[],
    241                 std::time::Instant::now() + Duration::from_millis(10),
    242             )
    243             .await
    244             .expect_err("empty relay list");
    245         assert_eq!(empty.code(), SafeErrorCode::InvalidRelayConfiguration);
    246 
    247         let mixed = [
    248             endpoint("ws://127.0.0.1:7777", RelayUrlPolicy::Local),
    249             endpoint("wss://relay.example", RelayUrlPolicy::Public),
    250         ];
    251         let error = adapter
    252             .fetch_profile(
    253                 public_key,
    254                 &mixed,
    255                 std::time::Instant::now() + Duration::from_millis(10),
    256             )
    257             .await
    258             .expect_err("mixed trust profiles");
    259         assert_eq!(error.code(), SafeErrorCode::InvalidRelayConfiguration);
    260     }
    261 
    262     #[tokio::test]
    263     async fn governed_transport_fails_when_no_relay_completes() {
    264         let error = SdkNostrClient::new(Duration::from_millis(25))
    265             .fetch_profile(
    266                 PublicKey::from_bytes([7; 32]).expect("public key"),
    267                 &[endpoint("ws://127.0.0.1:1", RelayUrlPolicy::Local)],
    268                 std::time::Instant::now() + Duration::from_millis(50),
    269             )
    270             .await
    271             .expect_err("unavailable relay");
    272         assert_eq!(error.code(), SafeErrorCode::RelayConnectionFailed);
    273     }
    274 
    275     #[test]
    276     fn debug_and_policy_are_safe_and_lib_owned() {
    277         let adapter = SdkNostrClient::new(Duration::from_secs(1));
    278         assert_eq!(
    279             format!("{adapter:?}"),
    280             "SdkNostrClient { timeout: 1s, transport: \"[sealed]\" }"
    281         );
    282         assert!(
    283             RelayEndpoint::new(
    284                 "wss://relay.example",
    285                 RelayUrlPolicy::Public,
    286                 RelayAccess::ReadOnly,
    287             )
    288             .is_ok()
    289         );
    290         assert!(
    291             RelayEndpoint::new(
    292                 "ws://relay.example",
    293                 RelayUrlPolicy::Public,
    294                 RelayAccess::ReadOnly,
    295             )
    296             .is_err()
    297         );
    298     }
    299 
    300     fn endpoint(value: &str, policy: RelayUrlPolicy) -> RelayEndpoint {
    301         RelayEndpoint::new(value, policy, RelayAccess::ReadWrite).expect("relay endpoint")
    302     }
    303 }