profile_refresh.rs (24395B)
1 use harvestcircle_domain::{PublicKey, SafeError, SafeErrorCode}; 2 use radroots_transport_nostr::RelayEndpoint; 3 use std::time::Instant; 4 5 use crate::{ 6 ActiveIdentitySnapshot, AppCore, AppSnapshot, CachedProfile, Clock, NostrClient, 7 ProfileFetchResult, ProfileLoadState, ProfileRefreshStatus, ProfileRepository, 8 RelayConnectionState, RelayFetchCompleteness, SnapshotRevision, StateTransition, 9 }; 10 11 #[derive(Clone, Debug, Eq, PartialEq)] 12 pub struct ProfileRefreshPlan { 13 public_key: PublicKey, 14 active_identity: ActiveIdentitySnapshot, 15 relays: Vec<RelayEndpoint>, 16 expected_revision: SnapshotRevision, 17 } 18 19 impl ProfileRefreshPlan { 20 #[must_use] 21 pub const fn public_key(&self) -> PublicKey { 22 self.public_key 23 } 24 25 #[must_use] 26 pub const fn active_identity(&self) -> &ActiveIdentitySnapshot { 27 &self.active_identity 28 } 29 30 #[must_use] 31 pub fn relays(&self) -> &[RelayEndpoint] { 32 &self.relays 33 } 34 35 #[must_use] 36 pub const fn expected_revision(&self) -> SnapshotRevision { 37 self.expected_revision 38 } 39 } 40 41 impl AppCore { 42 /// Manually refreshes the active identity's Nostr kind-0 profile. 43 /// 44 /// Cached public metadata remains visible while the asynchronous request is 45 /// running. Calling this command while signed out is an idempotent no-op. 46 /// 47 /// # Errors 48 /// 49 /// Returns a safe storage or application-state error. Relay and invalid-data 50 /// failures are represented as nonfatal snapshot state. 51 pub async fn refresh_active_profile( 52 &self, 53 profiles: &(impl ProfileRepository + ?Sized), 54 client: &(impl NostrClient + ?Sized), 55 clock: &(impl Clock + ?Sized), 56 deadline: Instant, 57 ) -> Result<AppSnapshot, SafeError> { 58 self.refresh_profile_for_active_identity(profiles, client, clock, deadline) 59 .await 60 } 61 62 /// Refreshes the current active identity while retaining any cached profile. 63 /// 64 /// Stale results are discarded when the identity is replaced or signed out. 65 /// 66 /// # Errors 67 /// 68 /// Returns a safe storage or application-state error. Relay and invalid-data 69 /// failures are represented as nonfatal snapshot state. 70 async fn refresh_profile_for_active_identity( 71 &self, 72 profiles: &(impl ProfileRepository + ?Sized), 73 client: &(impl NostrClient + ?Sized), 74 clock: &(impl Clock + ?Sized), 75 deadline: Instant, 76 ) -> Result<AppSnapshot, SafeError> { 77 let Some(plan) = self.begin_profile_refresh()? else { 78 return Ok(self.snapshot()); 79 }; 80 let result = client 81 .fetch_profile(plan.public_key(), plan.relays(), deadline) 82 .await; 83 self.complete_profile_refresh(&plan, result, profiles, clock) 84 .await 85 } 86 87 /// Begins a refresh on the actor and returns the immutable network plan. 88 /// 89 /// # Errors 90 /// 91 /// Returns a safe state error when the loading transition is invalid. 92 pub fn begin_profile_refresh(&self) -> Result<Option<ProfileRefreshPlan>, SafeError> { 93 let Some(active) = self.snapshot().active_identity().cloned() else { 94 return Ok(None); 95 }; 96 let public_key = active.identity().public_key(); 97 let loading = self.apply_transition(StateTransition::UpdateActiveIdentity { 98 expected: public_key, 99 active_identity: Box::new(ActiveIdentitySnapshot::new( 100 active.identity().clone(), 101 RelayConnectionState::Connecting, 102 ProfileLoadState::Loading, 103 active.profile().cloned(), 104 )), 105 problem: None, 106 })?; 107 Ok(Some(ProfileRefreshPlan { 108 public_key, 109 active_identity: active, 110 relays: loading.relay_configuration().relays().to_vec(), 111 expected_revision: loading.revision(), 112 })) 113 } 114 115 /// Applies a correlated refresh result on the actor. 116 /// 117 /// # Errors 118 /// 119 /// Returns a safe storage or application-state error. Stale results are 120 /// discarded without persistence or publication. 121 pub async fn complete_profile_refresh( 122 &self, 123 plan: &ProfileRefreshPlan, 124 result: Result<ProfileFetchResult, SafeError>, 125 profiles: &(impl ProfileRepository + ?Sized), 126 clock: &(impl Clock + ?Sized), 127 ) -> Result<AppSnapshot, SafeError> { 128 if !is_current_active(self, plan.public_key()) { 129 return Ok(self.snapshot()); 130 } 131 132 let current_active = self 133 .snapshot() 134 .active_identity() 135 .cloned() 136 .ok_or_else(invalid_profile_completion)?; 137 138 match result { 139 Ok(fetched) => { 140 let (candidate, completeness) = fetched.into_parts(); 141 self.complete_successful_profile_fetch( 142 plan, 143 current_active, 144 candidate, 145 completeness, 146 profiles, 147 clock, 148 ) 149 .await 150 } 151 Err(error) => { 152 let status = refresh_status(error); 153 profiles 154 .record_refresh_status(plan.public_key(), clock.now(), status) 155 .await?; 156 self.apply_transition(StateTransition::UpdateActiveIdentity { 157 expected: plan.public_key(), 158 active_identity: Box::new(ActiveIdentitySnapshot::new( 159 current_active.identity().clone(), 160 RelayConnectionState::Degraded, 161 ProfileLoadState::Error(error), 162 current_active.profile().cloned(), 163 )), 164 problem: Some(error), 165 }) 166 } 167 } 168 } 169 170 async fn complete_successful_profile_fetch( 171 &self, 172 plan: &ProfileRefreshPlan, 173 current_active: ActiveIdentitySnapshot, 174 candidate: Option<harvestcircle_domain::Kind0ProfileCandidate>, 175 completeness: RelayFetchCompleteness, 176 profiles: &(impl ProfileRepository + ?Sized), 177 clock: &(impl Clock + ?Sized), 178 ) -> Result<AppSnapshot, SafeError> { 179 let relay_state = match completeness { 180 RelayFetchCompleteness::Complete => RelayConnectionState::Connected, 181 RelayFetchCompleteness::Partial => RelayConnectionState::Degraded, 182 }; 183 let problem = match completeness { 184 RelayFetchCompleteness::Complete => None, 185 RelayFetchCompleteness::Partial => Some(partial_relay_result()), 186 }; 187 match candidate { 188 Some(candidate) => { 189 let cached = CachedProfile::new( 190 candidate.clone(), 191 clock.now(), 192 ProfileRefreshStatus::Success, 193 ); 194 profiles.save_profile(&cached).await?; 195 let winning_profile = profiles.load_profile(plan.public_key()).await?.map_or_else( 196 || candidate.metadata().clone(), 197 |profile| profile.candidate().metadata().clone(), 198 ); 199 self.apply_transition(StateTransition::UpdateActiveIdentity { 200 expected: plan.public_key(), 201 active_identity: Box::new(ActiveIdentitySnapshot::new( 202 current_active.identity().clone(), 203 relay_state, 204 ProfileLoadState::Fresh, 205 Some(winning_profile), 206 )), 207 problem, 208 }) 209 } 210 None => self.apply_transition(StateTransition::UpdateActiveIdentity { 211 expected: plan.public_key(), 212 active_identity: Box::new(ActiveIdentitySnapshot::new( 213 current_active.identity().clone(), 214 relay_state, 215 if current_active.profile().is_some() { 216 ProfileLoadState::Cached 217 } else { 218 ProfileLoadState::Empty 219 }, 220 current_active.profile().cloned(), 221 )), 222 problem, 223 }), 224 } 225 } 226 } 227 228 const fn partial_relay_result() -> SafeError { 229 SafeError::new( 230 SafeErrorCode::RelayConnectionFailed, 231 harvestcircle_domain::SafeMessage::new("One or more Nostr relays did not complete."), 232 ) 233 } 234 235 const fn invalid_profile_completion() -> SafeError { 236 SafeError::new( 237 SafeErrorCode::InvalidApplicationState, 238 harvestcircle_domain::SafeMessage::new("The active profile refresh is no longer valid."), 239 ) 240 } 241 242 fn is_current_active(core: &AppCore, public_key: PublicKey) -> bool { 243 core.snapshot() 244 .active_identity() 245 .is_some_and(|active| active.identity().public_key() == public_key) 246 } 247 248 const fn refresh_status(error: SafeError) -> ProfileRefreshStatus { 249 match error.code() { 250 SafeErrorCode::InvalidProfileMetadata | SafeErrorCode::ProfileRefreshFailed => { 251 ProfileRefreshStatus::InvalidData 252 } 253 _ => ProfileRefreshStatus::Offline, 254 } 255 } 256 257 #[cfg(test)] 258 mod tests { 259 use std::sync::Mutex; 260 use std::time::{Duration, Instant}; 261 262 use harvestcircle_domain::{ 263 EventId, Kind0ProfileCandidate, ProfileMetadata, PublicKey, SafeError, SafeErrorCode, 264 SafeMessage, SecretKeyInput, UnixTimestamp, select_latest_kind0, 265 }; 266 use radroots_transport_nostr::{RelayAccess, RelayEndpoint, RelayUrlPolicy}; 267 268 use crate::{ 269 ActiveIdentitySnapshot, AppCore, BoxFuture, CachedProfile, Clock, 270 InMemoryIdentityRepository, InMemoryOperationJournal, InMemorySecretStore, NostrClient, 271 ProfileFetchResult, ProfileLoadState, ProfileRefreshStatus, ProfileRepository, 272 RelayConfiguration, RelayConnectionState, 273 }; 274 275 #[derive(Default)] 276 struct MemoryProfiles(Mutex<Option<CachedProfile>>); 277 278 impl ProfileRepository for MemoryProfiles { 279 fn load_profile( 280 &self, 281 _public_key: PublicKey, 282 ) -> BoxFuture<'_, Result<Option<CachedProfile>, SafeError>> { 283 Box::pin(async move { Ok(self.0.lock().expect("profiles").clone()) }) 284 } 285 fn save_profile<'a>( 286 &'a self, 287 profile: &'a CachedProfile, 288 ) -> BoxFuture<'a, Result<(), SafeError>> { 289 Box::pin(async move { 290 let mut cached = self.0.lock().expect("profiles"); 291 let selected = cached.as_ref().map_or_else( 292 || profile.clone(), 293 |current| { 294 let winner = select_latest_kind0([ 295 current.candidate().clone(), 296 profile.candidate().clone(), 297 ]) 298 .expect("two candidates"); 299 if &winner == current.candidate() { 300 current.clone() 301 } else { 302 profile.clone() 303 } 304 }, 305 ); 306 *cached = Some(selected); 307 Ok(()) 308 }) 309 } 310 fn record_refresh_status<'a>( 311 &'a self, 312 _public_key: PublicKey, 313 refreshed_at: UnixTimestamp, 314 status: ProfileRefreshStatus, 315 ) -> BoxFuture<'a, Result<(), SafeError>> { 316 Box::pin(async move { 317 if let Some(profile) = self.0.lock().expect("profiles").as_mut() { 318 *profile = 319 CachedProfile::new(profile.candidate().clone(), refreshed_at, status); 320 } 321 Ok(()) 322 }) 323 } 324 fn remove_profile(&self, _public_key: PublicKey) -> BoxFuture<'_, Result<(), SafeError>> { 325 Box::pin(async move { 326 *self.0.lock().expect("profiles") = None; 327 Ok(()) 328 }) 329 } 330 } 331 332 struct FixedClock; 333 impl Clock for FixedClock { 334 fn now(&self) -> UnixTimestamp { 335 UnixTimestamp::from_seconds(50).expect("time") 336 } 337 } 338 339 struct FixedClient(Result<Option<Kind0ProfileCandidate>, SafeError>); 340 impl NostrClient for FixedClient { 341 fn fetch_profile<'a>( 342 &'a self, 343 _public_key: PublicKey, 344 _relays: &'a [RelayEndpoint], 345 _deadline: std::time::Instant, 346 ) -> BoxFuture<'a, Result<ProfileFetchResult, SafeError>> { 347 let result = self.0.clone(); 348 Box::pin(async move { result.map(ProfileFetchResult::complete) }) 349 } 350 } 351 352 struct BlockingClient { 353 started: tokio::sync::Semaphore, 354 release: tokio::sync::Semaphore, 355 result: Result<Option<Kind0ProfileCandidate>, SafeError>, 356 } 357 358 impl BlockingClient { 359 fn new(result: Result<Option<Kind0ProfileCandidate>, SafeError>) -> Self { 360 Self { 361 started: tokio::sync::Semaphore::new(0), 362 release: tokio::sync::Semaphore::new(0), 363 result, 364 } 365 } 366 } 367 368 impl NostrClient for BlockingClient { 369 fn fetch_profile<'a>( 370 &'a self, 371 _public_key: PublicKey, 372 _relays: &'a [RelayEndpoint], 373 _deadline: std::time::Instant, 374 ) -> BoxFuture<'a, Result<ProfileFetchResult, SafeError>> { 375 Box::pin(async move { 376 self.started.add_permits(1); 377 let permit = self.release.acquire().await.expect("release open"); 378 permit.forget(); 379 self.result.clone().map(ProfileFetchResult::complete) 380 }) 381 } 382 } 383 384 fn profile(public_key: PublicKey, name: &str, timestamp: i64) -> Kind0ProfileCandidate { 385 Kind0ProfileCandidate::new( 386 EventId::from_bytes([u8::try_from(timestamp).expect("small timestamp"); 32]), 387 public_key, 388 UnixTimestamp::from_seconds(timestamp).expect("time"), 389 ProfileMetadata::new(Some(name.to_owned()), None, None, None, None).expect("profile"), 390 ) 391 } 392 393 async fn active_core( 394 profiles: &MemoryProfiles, 395 cached_name: Option<&str>, 396 ) -> (AppCore, PublicKey) { 397 let relays = RelayConfiguration::new(vec![ 398 RelayEndpoint::new( 399 "ws://localhost:8080", 400 RelayUrlPolicy::Local, 401 RelayAccess::ReadWrite, 402 ) 403 .expect("relay"), 404 ]) 405 .expect("relay configuration"); 406 let core = AppCore::in_memory(relays); 407 let identities = InMemoryIdentityRepository::default(); 408 let secrets = InMemorySecretStore::default(); 409 let journal = InMemoryOperationJournal::default(); 410 core.bootstrap().expect("bootstrap"); 411 let public_key = core 412 .import_secret_key( 413 SecretKeyInput::parse( 414 "7e7e9c42a91bfef19fa7ea99d52d8afdb67d893a8fefba1f5cb9793f2107f6d7".to_owned(), 415 ) 416 .expect("secret"), 417 &identities, 418 &identities, 419 &secrets, 420 &journal, 421 &FixedClock, 422 ) 423 .await 424 .expect("import") 425 .identity() 426 .public_key(); 427 if let Some(name) = cached_name { 428 profiles 429 .save_profile(&CachedProfile::new( 430 profile(public_key, name, 10), 431 UnixTimestamp::from_seconds(11).expect("time"), 432 ProfileRefreshStatus::Success, 433 )) 434 .await 435 .expect("cache"); 436 } 437 core.activate_identity( 438 public_key, 439 &identities, 440 &identities, 441 profiles, 442 &secrets, 443 &FixedClock, 444 ) 445 .await 446 .expect("activate"); 447 (core, public_key) 448 } 449 450 #[tokio::test] 451 async fn refresh_transitions_from_cache_through_loading_to_fresh_profile() { 452 let profiles = MemoryProfiles::default(); 453 let (core, public_key) = active_core(&profiles, Some("Cached")).await; 454 assert_eq!( 455 core.snapshot() 456 .active_identity() 457 .map(crate::ActiveIdentitySnapshot::profile_state), 458 Some(ProfileLoadState::Cached) 459 ); 460 let plan = core 461 .begin_profile_refresh() 462 .expect("begin refresh") 463 .expect("active refresh"); 464 let loading = core.snapshot(); 465 assert_eq!( 466 loading 467 .active_identity() 468 .map(crate::ActiveIdentitySnapshot::profile_state), 469 Some(ProfileLoadState::Loading) 470 ); 471 assert_eq!( 472 loading 473 .active_identity() 474 .map(crate::ActiveIdentitySnapshot::relay_state), 475 Some(RelayConnectionState::Connecting) 476 ); 477 let client = FixedClient(Ok(Some(profile(public_key, "Fresh", 20)))); 478 let result = client 479 .fetch_profile(plan.public_key(), plan.relays(), deadline()) 480 .await; 481 core.complete_profile_refresh(&plan, result, &profiles, &FixedClock) 482 .await 483 .expect("complete refresh"); 484 assert_eq!( 485 core.snapshot() 486 .active_identity() 487 .and_then(|active| active.profile()) 488 .and_then(ProfileMetadata::name), 489 Some("Fresh") 490 ); 491 } 492 493 #[tokio::test] 494 async fn refresh_failure_preserves_cached_profile_as_nonfatal_state() { 495 let profiles = MemoryProfiles::default(); 496 let (core, public_key) = active_core(&profiles, Some("Cached")).await; 497 let cached = profile(public_key, "Cached", 10); 498 let error = SafeError::new( 499 SafeErrorCode::RelayConnectionFailed, 500 SafeMessage::new("The relay is offline."), 501 ); 502 503 let snapshot = core 504 .refresh_profile_for_active_identity( 505 &profiles, 506 &FixedClient(Err(error)), 507 &FixedClock, 508 deadline(), 509 ) 510 .await 511 .expect("nonfatal refresh"); 512 513 assert_eq!(snapshot.recoverable_problem(), Some(error)); 514 assert_eq!( 515 snapshot 516 .active_identity() 517 .map(crate::ActiveIdentitySnapshot::relay_state), 518 Some(RelayConnectionState::Degraded) 519 ); 520 assert_eq!( 521 profiles 522 .load_profile(public_key) 523 .await 524 .expect("load") 525 .expect("cache") 526 .candidate(), 527 &cached 528 ); 529 } 530 531 #[tokio::test] 532 async fn refresh_discards_stale_completion_after_sign_out() { 533 let profiles = MemoryProfiles::default(); 534 let (core, public_key) = active_core(&profiles, Some("Cached")).await; 535 let client = BlockingClient::new(Ok(Some(profile(public_key, "Stale", 20)))); 536 537 let refresh = 538 core.refresh_profile_for_active_identity(&profiles, &client, &FixedClock, deadline()); 539 let sign_out = async { 540 let permit = client.started.acquire().await.expect("refresh starts"); 541 permit.forget(); 542 core.sign_out().expect("sign out"); 543 client.release.add_permits(1); 544 }; 545 let (result, ()) = tokio::join!(refresh, sign_out); 546 547 assert!( 548 result 549 .expect("stale result is harmless") 550 .active_identity() 551 .is_none() 552 ); 553 assert_eq!( 554 profiles 555 .load_profile(public_key) 556 .await 557 .expect("load") 558 .expect("cached") 559 .candidate() 560 .metadata() 561 .name(), 562 Some("Cached") 563 ); 564 } 565 566 #[tokio::test] 567 async fn manual_refresh_is_repeatable_and_signed_out_safe() { 568 let profiles = MemoryProfiles::default(); 569 let (core, public_key) = active_core(&profiles, None).await; 570 let first = core 571 .refresh_active_profile( 572 &profiles, 573 &FixedClient(Ok(Some(profile(public_key, "First", 10)))), 574 &FixedClock, 575 deadline(), 576 ) 577 .await 578 .expect("first refresh"); 579 let second = core 580 .refresh_active_profile( 581 &profiles, 582 &FixedClient(Ok(Some(profile(public_key, "Second", 20)))), 583 &FixedClock, 584 deadline(), 585 ) 586 .await 587 .expect("second refresh"); 588 589 assert!(second.revision() > first.revision()); 590 assert_eq!( 591 second 592 .active_identity() 593 .and_then(|active| active.profile()) 594 .and_then(ProfileMetadata::name), 595 Some("Second") 596 ); 597 let signed_out = core.sign_out().expect("sign out"); 598 let no_op = core 599 .refresh_active_profile(&profiles, &FixedClient(Ok(None)), &FixedClock, deadline()) 600 .await 601 .expect("signed-out no-op"); 602 assert_eq!(no_op, signed_out); 603 } 604 605 fn deadline() -> Instant { 606 Instant::now() + Duration::from_secs(5) 607 } 608 609 #[tokio::test] 610 async fn partial_relay_success_retains_verified_data_and_marks_degraded_connectivity() { 611 let profiles = MemoryProfiles::default(); 612 let (core, public_key) = active_core(&profiles, None).await; 613 let plan = core 614 .begin_profile_refresh() 615 .expect("begin") 616 .expect("active"); 617 let snapshot = core 618 .complete_profile_refresh( 619 &plan, 620 Ok(ProfileFetchResult::partial(Some(profile( 621 public_key, "Partial", 30, 622 )))), 623 &profiles, 624 &FixedClock, 625 ) 626 .await 627 .expect("partial completion"); 628 let active = snapshot.active_identity().expect("active identity"); 629 assert_eq!(active.relay_state(), RelayConnectionState::Degraded); 630 assert_eq!(active.profile_state(), ProfileLoadState::Fresh); 631 assert_eq!( 632 active.profile().and_then(ProfileMetadata::name), 633 Some("Partial") 634 ); 635 assert_eq!( 636 snapshot.recoverable_problem().map(SafeError::code), 637 Some(SafeErrorCode::RelayConnectionFailed) 638 ); 639 } 640 641 #[tokio::test] 642 async fn overlapping_refreshes_keep_the_newest_event_regardless_of_completion_order() { 643 let profiles = MemoryProfiles::default(); 644 let (core, public_key) = active_core(&profiles, Some("Cached")).await; 645 let first = core 646 .begin_profile_refresh() 647 .expect("first") 648 .expect("active"); 649 let second = core 650 .begin_profile_refresh() 651 .expect("second") 652 .expect("active"); 653 654 core.complete_profile_refresh( 655 &second, 656 Ok(ProfileFetchResult::complete(Some(profile( 657 public_key, "Newest", 30, 658 )))), 659 &profiles, 660 &FixedClock, 661 ) 662 .await 663 .expect("newest completes first"); 664 let final_snapshot = core 665 .complete_profile_refresh( 666 &first, 667 Ok(ProfileFetchResult::complete(Some(profile( 668 public_key, "Older", 20, 669 )))), 670 &profiles, 671 &FixedClock, 672 ) 673 .await 674 .expect("older completes last"); 675 676 assert_eq!( 677 final_snapshot 678 .active_identity() 679 .and_then(ActiveIdentitySnapshot::profile) 680 .and_then(ProfileMetadata::name), 681 Some("Newest") 682 ); 683 assert_eq!( 684 profiles 685 .load_profile(public_key) 686 .await 687 .expect("cache") 688 .expect("profile") 689 .candidate() 690 .metadata() 691 .name(), 692 Some("Newest") 693 ); 694 } 695 }