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 }