runtime_actor.rs (94613B)
1 use std::collections::BTreeMap; 2 use std::future::Future; 3 use std::num::{NonZeroU64, NonZeroUsize}; 4 use std::sync::atomic::{AtomicU64, Ordering}; 5 use std::sync::{Arc, Mutex}; 6 use std::time::{Duration, Instant}; 7 8 use harvestcircle_application::{ 9 ActiveSessionBinding, ActorMailbox, AppSnapshot, ChangeSubscriptionId, Clock, CommandContext, 10 CommandEnvelope, CommandReceipt, CommandResult, CommandSubmission, DurableRequestId, 11 GenerateIdentityReceipt, GeneratedKeyRecoveryHandle, GeneratedKeyStage, ImportIdentityReceipt, 12 LifecycleGate, NostrClient, OrderedSnapshotChanges, ProfileFetchResult, ProfileRefreshPlan, 13 RecoveryStageId, RelayConfiguration, RemovalConfirmationToken, RequestId, RuntimeCommandClass, 14 RuntimeLifecycle, SecretStore, SessionGeneration, SnapshotChange, SnapshotChangeReceiver, 15 SnapshotRevision, StagedGeneratedKey, TaskCorrelation, 16 }; 17 use harvestcircle_domain::{ 18 LocalKeyringBinding, NostrIdentityReference, PublicKey, SafeError, SafeErrorCode, SafeMessage, 19 SecretKeyInput, SignerAvailability, 20 }; 21 use radroots_runtime_paths::RuntimeContext; 22 #[cfg(test)] 23 use radroots_runtime_paths::{ 24 InstanceId, RadrootsHostEnvironment, RadrootsPathProfile, RadrootsPathResolver, 25 RadrootsPlatform, RuntimeContextBootstrap, RuntimeContextSource, ServiceId, 26 }; 27 use radroots_service_sqlite::MigrationBuildIdentity; 28 use tokio::runtime::Handle; 29 use tokio::sync::{mpsc, oneshot, watch}; 30 31 use crate::{InstallationIdentity, InstallationIdentitySource, PersistentAppCore}; 32 33 const DEFAULT_COMMAND_TIMEOUT: Duration = Duration::from_secs(30); 34 const DEFAULT_TASK_CAPACITY: usize = 64; 35 36 enum RuntimeCommand { 37 Snapshot, 38 GenerateIdentity { 39 durable_request: DurableRequestId, 40 expected_revision: u64, 41 }, 42 BeginGeneratedKeyStage, 43 AcknowledgeGeneratedKeyStage { 44 id: RecoveryStageId, 45 durable_request: DurableRequestId, 46 }, 47 CancelGeneratedKeyStage, 48 ImportSecretKey { 49 input: SecretKeyInput, 50 durable_request: DurableRequestId, 51 expected_revision: u64, 52 }, 53 SelectIdentity(PublicKey), 54 ActivateIdentity(PublicKey), 55 SignOut, 56 RefreshActiveProfile, 57 RequestIdentityRemoval(PublicKey), 58 ConfirmIdentityRemoval { 59 token: RemovalConfirmationToken, 60 durable_request: DurableRequestId, 61 }, 62 SubscribeChanges(NonZeroUsize), 63 UnsubscribeChanges(ChangeSubscriptionId), 64 Close, 65 } 66 67 enum RuntimeCommandValue { 68 Snapshot(Box<AppSnapshot>), 69 Generated(GenerateIdentityReceipt), 70 GeneratedKeyStage(GeneratedKeyRecoveryHandle), 71 GeneratedKeyStageCancelled(bool), 72 Imported(ImportIdentityReceipt), 73 RemovalRequest(RemovalConfirmationToken), 74 Subscription(RuntimeChangeSubscription), 75 Unsubscribed(bool), 76 Closed, 77 } 78 79 impl RuntimeCommand { 80 const fn class(&self) -> RuntimeCommandClass { 81 match self { 82 Self::Snapshot | Self::SubscribeChanges(_) | Self::UnsubscribeChanges(_) => { 83 RuntimeCommandClass::Observe 84 } 85 Self::GenerateIdentity { .. } 86 | Self::BeginGeneratedKeyStage 87 | Self::AcknowledgeGeneratedKeyStage { .. } 88 | Self::ImportSecretKey { .. } 89 | Self::ActivateIdentity(_) 90 | Self::ConfirmIdentityRemoval { .. } => RuntimeCommandClass::UseCredential, 91 Self::SelectIdentity(_) 92 | Self::SignOut 93 | Self::RequestIdentityRemoval(_) 94 | Self::CancelGeneratedKeyStage => RuntimeCommandClass::MutateLocalState, 95 Self::RefreshActiveProfile => RuntimeCommandClass::UseRelay, 96 Self::Close => RuntimeCommandClass::Shutdown, 97 } 98 } 99 100 const fn resolves_revision_through_durable_replay(&self) -> bool { 101 matches!( 102 self, 103 Self::GenerateIdentity { .. } | Self::ImportSecretKey { .. } 104 ) 105 } 106 } 107 108 struct RuntimeActor { 109 adapter: Arc<PersistentAppCore>, 110 secrets: Arc<dyn SecretStore>, 111 clock: Arc<dyn Clock>, 112 nostr: Arc<dyn NostrClient>, 113 lifecycle: Arc<Mutex<LifecycleGate>>, 114 runtime: Handle, 115 session_generation: SessionGeneration, 116 published_session_generation: Arc<AtomicU64>, 117 profile_tasks: BTreeMap<RequestId, PendingProfileTask>, 118 changes: OrderedSnapshotChanges, 119 published_foreground_session: Arc<Mutex<Option<ActiveSessionBinding>>>, 120 generated_key_stage: GeneratedKeyStage, 121 } 122 123 struct PendingProfileTask { 124 correlation: TaskCorrelation, 125 plan: ProfileRefreshPlan, 126 deadline: Instant, 127 reply: oneshot::Sender<CommandReceipt<RuntimeCommandValue>>, 128 handle: tokio::task::JoinHandle<()>, 129 } 130 131 struct ProfileCompletion { 132 request_id: RequestId, 133 result: Result<ProfileFetchResult, SafeError>, 134 } 135 136 #[derive(Clone)] 137 pub struct RuntimeActorHandle { 138 mailbox: ActorMailbox<RuntimeCommand, RuntimeCommandValue>, 139 adapter: Arc<PersistentAppCore>, 140 lifecycle: Arc<Mutex<LifecycleGate>>, 141 next_request: Arc<AtomicU64>, 142 session_generation: Arc<AtomicU64>, 143 foreground_session: Arc<Mutex<Option<ActiveSessionBinding>>>, 144 installation_identity: InstallationIdentity, 145 runtime: Handle, 146 actor_task: Arc<Mutex<Option<tokio::task::JoinHandle<()>>>>, 147 actor_exit: watch::Receiver<bool>, 148 #[cfg(test)] 149 _test_directory: Option<Arc<tempfile::TempDir>>, 150 } 151 152 #[derive(Clone)] 153 pub struct RuntimeDependencies { 154 secrets: Arc<dyn SecretStore>, 155 clock: Arc<dyn Clock>, 156 nostr: Arc<dyn NostrClient>, 157 installation_source: Arc<dyn InstallationIdentitySource>, 158 } 159 160 impl RuntimeDependencies { 161 #[must_use] 162 pub fn new( 163 secrets: Arc<dyn SecretStore>, 164 clock: Arc<dyn Clock>, 165 nostr: Arc<dyn NostrClient>, 166 installation_source: Arc<dyn InstallationIdentitySource>, 167 ) -> Self { 168 Self { 169 secrets, 170 clock, 171 nostr, 172 installation_source, 173 } 174 } 175 } 176 177 pub struct RuntimeChangeSubscription { 178 id: ChangeSubscriptionId, 179 receiver: SnapshotChangeReceiver, 180 } 181 182 impl RuntimeChangeSubscription { 183 #[must_use] 184 pub const fn id(&self) -> ChangeSubscriptionId { 185 self.id 186 } 187 188 pub async fn receive(&mut self) -> Option<SnapshotChange> { 189 self.receiver.receive().await 190 } 191 } 192 193 impl RuntimeActorHandle { 194 /// Opens, migrates, recovers, and starts one actor-owned file-backed runtime. 195 /// 196 /// # Errors 197 /// 198 /// Returns a safe storage, recovery, or lifecycle error before the actor is 199 /// published when opening cannot reach ready state. 200 pub async fn open( 201 context: &RuntimeContext, 202 relay_configuration: RelayConfiguration, 203 dependencies: RuntimeDependencies, 204 build: &MigrationBuildIdentity, 205 capacity: NonZeroUsize, 206 runtime: &Handle, 207 ) -> Result<Self, SafeError> { 208 let applied_at_unix_s = u64::try_from(dependencies.clock.now().as_seconds()) 209 .map_err(|_| invalid_runtime_evidence())?; 210 let created_at_unix_ms = applied_at_unix_s 211 .checked_mul(1_000) 212 .ok_or_else(invalid_runtime_evidence)?; 213 let adapter = PersistentAppCore::open( 214 context, 215 relay_configuration, 216 created_at_unix_ms, 217 applied_at_unix_s, 218 build, 219 ) 220 .await?; 221 Self::start(adapter, dependencies, capacity, runtime).await 222 } 223 224 async fn start( 225 adapter: PersistentAppCore, 226 dependencies: RuntimeDependencies, 227 capacity: NonZeroUsize, 228 runtime: &Handle, 229 ) -> Result<Self, SafeError> { 230 let mut gate = LifecycleGate::opening(); 231 gate.begin_compatibility_check()?; 232 gate.compatibility_accepted()?; 233 gate.ownership_acquired()?; 234 gate.migration_complete()?; 235 let RuntimeDependencies { 236 secrets, 237 clock, 238 nostr, 239 installation_source, 240 } = dependencies; 241 let adapter = Arc::new(adapter); 242 adapter.bootstrap(secrets.as_ref(), clock.as_ref()).await?; 243 let installation_identity = adapter 244 .initialize_installation_identity(installation_source.as_ref()) 245 .await?; 246 gate.recovery_complete()?; 247 248 let lifecycle = Arc::new(Mutex::new(gate)); 249 let (mailbox, receiver) = ActorMailbox::bounded(capacity); 250 let session_generation = Arc::new(AtomicU64::new(SessionGeneration::initial().value())); 251 let foreground_session = Arc::new(Mutex::new(None)); 252 let changes = OrderedSnapshotChanges::new(adapter.core().snapshot()); 253 let actor = RuntimeActor { 254 adapter: Arc::clone(&adapter), 255 secrets, 256 clock, 257 nostr, 258 lifecycle: Arc::clone(&lifecycle), 259 runtime: runtime.clone(), 260 session_generation: SessionGeneration::initial(), 261 published_session_generation: Arc::clone(&session_generation), 262 profile_tasks: BTreeMap::new(), 263 changes, 264 published_foreground_session: Arc::clone(&foreground_session), 265 generated_key_stage: GeneratedKeyStage::default(), 266 }; 267 let (actor_exit_sender, actor_exit) = watch::channel(false); 268 let actor_task = runtime.spawn(async move { 269 actor.run(receiver).await; 270 let _ = actor_exit_sender.send(true); 271 }); 272 Ok(Self { 273 mailbox, 274 adapter, 275 lifecycle, 276 next_request: Arc::new(AtomicU64::new(1)), 277 session_generation, 278 foreground_session, 279 installation_identity, 280 runtime: runtime.clone(), 281 actor_task: Arc::new(Mutex::new(Some(actor_task))), 282 actor_exit, 283 #[cfg(test)] 284 _test_directory: None, 285 }) 286 } 287 288 #[cfg(test)] 289 async fn in_memory( 290 relay_configuration: RelayConfiguration, 291 dependencies: RuntimeDependencies, 292 capacity: NonZeroUsize, 293 runtime: &Handle, 294 ) -> Result<Self, SafeError> { 295 let directory = Arc::new(tempfile::tempdir().map_err(|_| invalid_runtime_evidence())?); 296 let context = test_runtime_context(directory.path())?; 297 let build = test_migration_build_identity()?; 298 let mut actor = Self::open( 299 &context, 300 relay_configuration, 301 dependencies, 302 &build, 303 capacity, 304 runtime, 305 ) 306 .await?; 307 actor._test_directory = Some(directory); 308 Ok(actor) 309 } 310 311 #[must_use] 312 pub fn lifecycle(&self) -> RuntimeLifecycle { 313 self.lifecycle.lock().map_or_else( 314 |_| RuntimeLifecycle::Fatal(runtime_state_unavailable()), 315 |lifecycle| lifecycle.lifecycle(), 316 ) 317 } 318 319 #[must_use] 320 pub fn session_generation(&self) -> SessionGeneration { 321 SessionGeneration::from_value(self.session_generation.load(Ordering::Acquire)) 322 } 323 324 #[must_use] 325 pub fn foreground_session(&self) -> Option<ActiveSessionBinding> { 326 self.foreground_session 327 .lock() 328 .map(|session| session.clone()) 329 .unwrap_or(None) 330 } 331 332 #[must_use] 333 pub fn installation_identity(&self) -> &InstallationIdentity { 334 &self.installation_identity 335 } 336 337 #[must_use] 338 pub fn snapshot(&self) -> AppSnapshot { 339 self.adapter.core().snapshot() 340 } 341 342 /// Returns the ready snapshot through the actor command boundary. 343 /// 344 /// # Errors 345 /// 346 /// Returns a typed safe actor error. 347 pub async fn bootstrap(&self) -> Result<AppSnapshot, SafeError> { 348 Self::expect_snapshot(self.dispatch(RuntimeCommand::Snapshot, None).await?) 349 } 350 351 /// Generates one identity through the serialized actor boundary. 352 /// 353 /// # Errors 354 /// 355 /// Returns a safe identity, storage, keyring, timeout, or actor error. 356 pub async fn generate_identity( 357 &self, 358 request: DurableRequestId, 359 expected_revision: SnapshotRevision, 360 timeout: Duration, 361 ) -> Result<GenerateIdentityReceipt, SafeError> { 362 match self 363 .dispatch_durable( 364 RuntimeCommand::GenerateIdentity { 365 durable_request: request, 366 expected_revision: expected_revision.value(), 367 }, 368 expected_revision, 369 timeout, 370 ) 371 .await? 372 { 373 RuntimeCommandValue::Generated(receipt) => Ok(receipt), 374 _ => Err(invalid_actor_response()), 375 } 376 } 377 378 /// Begins the only actor-owned generated-key recovery stage. 379 /// 380 /// # Errors 381 /// 382 /// Returns a safe conflict, timeout, key-generation, or actor error. 383 pub async fn begin_generated_key_stage(&self) -> Result<GeneratedKeyRecoveryHandle, SafeError> { 384 match self 385 .dispatch(RuntimeCommand::BeginGeneratedKeyStage, None) 386 .await? 387 { 388 RuntimeCommandValue::GeneratedKeyStage(view) => Ok(view), 389 _ => Err(invalid_actor_response()), 390 } 391 } 392 393 /// Acknowledges recovery and commits the staged identity and credential once. 394 /// 395 /// # Errors 396 /// 397 /// Returns a safe unavailable, conflict, keyring, storage, timeout, or actor error. 398 pub async fn acknowledge_generated_key_stage( 399 &self, 400 id: RecoveryStageId, 401 request: DurableRequestId, 402 expected_revision: SnapshotRevision, 403 timeout: Duration, 404 ) -> Result<AppSnapshot, SafeError> { 405 let value = self 406 .dispatch_durable( 407 RuntimeCommand::AcknowledgeGeneratedKeyStage { 408 id, 409 durable_request: request, 410 }, 411 expected_revision, 412 timeout, 413 ) 414 .await?; 415 Self::expect_snapshot(value) 416 } 417 418 /// Cancels and zeroizes the active generated-key stage, if present. 419 /// 420 /// # Errors 421 /// 422 /// Returns a safe timeout or actor error. 423 pub async fn cancel_generated_key_stage(&self) -> Result<bool, SafeError> { 424 match self 425 .dispatch(RuntimeCommand::CancelGeneratedKeyStage, None) 426 .await? 427 { 428 RuntimeCommandValue::GeneratedKeyStageCancelled(cancelled) => Ok(cancelled), 429 _ => Err(invalid_actor_response()), 430 } 431 } 432 433 /// Imports one identity through the serialized actor boundary. 434 /// 435 /// # Errors 436 /// 437 /// Returns a safe identity, storage, keyring, timeout, or actor error. 438 pub async fn import_secret_key( 439 &self, 440 request: DurableRequestId, 441 expected_revision: SnapshotRevision, 442 input: SecretKeyInput, 443 timeout: Duration, 444 ) -> Result<ImportIdentityReceipt, SafeError> { 445 match self 446 .dispatch_durable( 447 RuntimeCommand::ImportSecretKey { 448 input, 449 durable_request: request, 450 expected_revision: expected_revision.value(), 451 }, 452 expected_revision, 453 timeout, 454 ) 455 .await? 456 { 457 RuntimeCommandValue::Imported(receipt) => Ok(receipt), 458 _ => Err(invalid_actor_response()), 459 } 460 } 461 462 /// Selects one identity through the serialized actor boundary. 463 /// 464 /// # Errors 465 /// 466 /// Returns a safe identity, storage, timeout, or actor error. 467 pub async fn select_identity(&self, public_key: PublicKey) -> Result<AppSnapshot, SafeError> { 468 let value = self 469 .dispatch(RuntimeCommand::SelectIdentity(public_key), None) 470 .await?; 471 Self::expect_snapshot(value) 472 } 473 474 /// Activates one identity through the serialized actor boundary. 475 /// 476 /// # Errors 477 /// 478 /// Returns a safe identity, credential, storage, timeout, or actor error. 479 pub async fn activate_identity(&self, public_key: PublicKey) -> Result<AppSnapshot, SafeError> { 480 let value = self 481 .dispatch(RuntimeCommand::ActivateIdentity(public_key), None) 482 .await?; 483 Self::expect_snapshot(value) 484 } 485 486 /// Signs out through the serialized actor boundary. 487 /// 488 /// # Errors 489 /// 490 /// Returns a safe timeout or actor error. 491 pub async fn sign_out(&self) -> Result<AppSnapshot, SafeError> { 492 let value = self.dispatch(RuntimeCommand::SignOut, None).await?; 493 Self::expect_snapshot(value) 494 } 495 496 /// Refreshes the active profile through the serialized actor boundary. 497 /// 498 /// # Errors 499 /// 500 /// Returns a safe relay, storage, timeout, or actor error. 501 pub async fn refresh_active_profile(&self) -> Result<AppSnapshot, SafeError> { 502 let value = self 503 .dispatch(RuntimeCommand::RefreshActiveProfile, None) 504 .await?; 505 Self::expect_snapshot(value) 506 } 507 508 /// Creates one removal request through the serialized actor boundary. 509 /// 510 /// # Errors 511 /// 512 /// Returns a safe identity, timeout, or actor error. 513 pub async fn request_identity_removal( 514 &self, 515 public_key: PublicKey, 516 ) -> Result<RemovalConfirmationToken, SafeError> { 517 match self 518 .dispatch(RuntimeCommand::RequestIdentityRemoval(public_key), None) 519 .await? 520 { 521 RuntimeCommandValue::RemovalRequest(token) => Ok(token), 522 _ => Err(invalid_actor_response()), 523 } 524 } 525 526 /// Confirms one removal through the serialized actor boundary. 527 /// 528 /// # Errors 529 /// 530 /// Returns a safe identity, credential, storage, timeout, or actor error. 531 pub async fn confirm_identity_removal( 532 &self, 533 token: RemovalConfirmationToken, 534 request: DurableRequestId, 535 expected_revision: SnapshotRevision, 536 timeout: Duration, 537 ) -> Result<AppSnapshot, SafeError> { 538 let value = self 539 .dispatch_durable( 540 RuntimeCommand::ConfirmIdentityRemoval { 541 token, 542 durable_request: request, 543 }, 544 expected_revision, 545 timeout, 546 ) 547 .await?; 548 Self::expect_snapshot(value) 549 } 550 551 /// Closes command admission and cancels supervised work. 552 /// 553 /// # Errors 554 /// 555 /// Returns a safe timeout or actor error. Repeated calls return closed. 556 pub async fn close(&self) -> Result<(), SafeError> { 557 self.close_with_timeout(DEFAULT_COMMAND_TIMEOUT).await 558 } 559 560 /// Closes the runtime within the supplied command deadline. 561 /// 562 /// # Errors 563 /// 564 /// Returns a safe timeout or actor error. An expired queued close cannot 565 /// later change runtime state. 566 pub async fn close_with_timeout(&self, timeout: Duration) -> Result<(), SafeError> { 567 let deadline = Instant::now() + timeout; 568 if matches!(self.lifecycle(), RuntimeLifecycle::Closed) { 569 return self.await_actor_exit(deadline).await; 570 } 571 let raw_request = self.next_request.fetch_add(1, Ordering::Relaxed); 572 let request_id = RequestId::new(raw_request).ok_or_else(request_space_exhausted)?; 573 let command_result = match self 574 .dispatch_with_deadline(RuntimeCommand::Close, None, request_id, deadline) 575 .await 576 { 577 Ok(RuntimeCommandValue::Closed) => Ok(()), 578 Ok(_) => Err(invalid_actor_response()), 579 Err(error) => Err(error), 580 }; 581 let exit_result = self.await_actor_exit(deadline).await; 582 if matches!(self.lifecycle(), RuntimeLifecycle::Closed) { 583 exit_result 584 } else { 585 command_result.and(exit_result) 586 } 587 } 588 589 async fn await_actor_exit(&self, deadline: Instant) -> Result<(), SafeError> { 590 let mut actor_exit = self.actor_exit.clone(); 591 if !*actor_exit.borrow() { 592 let remaining = deadline.saturating_duration_since(Instant::now()); 593 self.timeout(remaining, actor_exit.wait_for(|exited| *exited)) 594 .await 595 .map_err(|_| command_timed_out())? 596 .map_err(|_| runtime_closed())?; 597 } 598 let actor_task = self 599 .actor_task 600 .lock() 601 .map_err(|_| runtime_state_unavailable())? 602 .take(); 603 if let Some(actor_task) = actor_task { 604 let remaining = deadline.saturating_duration_since(Instant::now()); 605 self.timeout(remaining, actor_task) 606 .await 607 .map_err(|_| command_timed_out())? 608 .map_err(|_| runtime_closed())?; 609 } 610 Ok(()) 611 } 612 613 /// Atomically registers a bounded ordered change consumer with its initial snapshot. 614 /// 615 /// # Errors 616 /// 617 /// Returns a safe actor or subscription error. 618 pub async fn subscribe_changes( 619 &self, 620 capacity: NonZeroUsize, 621 ) -> Result<RuntimeChangeSubscription, SafeError> { 622 match self 623 .dispatch(RuntimeCommand::SubscribeChanges(capacity), None) 624 .await? 625 { 626 RuntimeCommandValue::Subscription(subscription) => Ok(subscription), 627 _ => Err(invalid_actor_response()), 628 } 629 } 630 631 /// Removes a change consumer through the serialized actor boundary. 632 /// 633 /// # Errors 634 /// 635 /// Returns a safe actor error. 636 pub async fn unsubscribe_changes(&self, id: ChangeSubscriptionId) -> Result<bool, SafeError> { 637 match self 638 .dispatch(RuntimeCommand::UnsubscribeChanges(id), None) 639 .await? 640 { 641 RuntimeCommandValue::Unsubscribed(removed) => Ok(removed), 642 _ => Err(invalid_actor_response()), 643 } 644 } 645 646 async fn dispatch( 647 &self, 648 command: RuntimeCommand, 649 expected_revision: Option<SnapshotRevision>, 650 ) -> Result<RuntimeCommandValue, SafeError> { 651 let raw_request = self.next_request.fetch_add(1, Ordering::Relaxed); 652 let request_id = RequestId::new(raw_request).ok_or_else(request_space_exhausted)?; 653 self.dispatch_with_deadline( 654 command, 655 expected_revision, 656 request_id, 657 Instant::now() + DEFAULT_COMMAND_TIMEOUT, 658 ) 659 .await 660 } 661 662 async fn dispatch_durable( 663 &self, 664 command: RuntimeCommand, 665 expected_revision: SnapshotRevision, 666 timeout: Duration, 667 ) -> Result<RuntimeCommandValue, SafeError> { 668 let raw_request = self.next_request.fetch_add(1, Ordering::Relaxed); 669 let request_id = RequestId::new(raw_request).ok_or_else(request_space_exhausted)?; 670 self.dispatch_with_deadline( 671 command, 672 Some(expected_revision), 673 request_id, 674 Instant::now() + timeout, 675 ) 676 .await 677 } 678 679 async fn dispatch_with_deadline( 680 &self, 681 command: RuntimeCommand, 682 expected_revision: Option<SnapshotRevision>, 683 request_id: RequestId, 684 deadline: Instant, 685 ) -> Result<RuntimeCommandValue, SafeError> { 686 let context = CommandContext::new(request_id, expected_revision, deadline); 687 let receipt = match self.mailbox.submit(context, command) { 688 CommandSubmission::Accepted(ticket) => { 689 let remaining = deadline.saturating_duration_since(Instant::now()); 690 match self.timeout(remaining, ticket.receipt()).await { 691 Ok(receipt) => receipt, 692 Err(_) => CommandReceipt::new(request_id, CommandResult::TimedOut), 693 } 694 } 695 CommandSubmission::Rejected(receipt) => receipt, 696 }; 697 match receipt.into_result() { 698 CommandResult::Completed(value) => Ok(value), 699 CommandResult::Conflicted { .. } => Err(command_conflicted()), 700 CommandResult::Rejected(_) => Err(command_rejected()), 701 CommandResult::TimedOut => Err(command_timed_out()), 702 CommandResult::Closed => Err(runtime_closed()), 703 CommandResult::Failed(error) => Err(error), 704 } 705 } 706 707 fn timeout<F>(&self, duration: Duration, future: F) -> tokio::time::Timeout<F> 708 where 709 F: Future, 710 { 711 let _guard = self.runtime.enter(); 712 tokio::time::timeout(duration, future) 713 } 714 715 #[cfg(test)] 716 async fn import_secret_key_test( 717 &self, 718 input: SecretKeyInput, 719 ) -> Result<ImportIdentityReceipt, SafeError> { 720 let request_number = self.next_request.fetch_add(1, Ordering::Relaxed); 721 self.import_secret_key( 722 DurableRequestId::parse(format!( 723 "01890f3e-7b1c-7000-8000-{:012x}", 724 request_number + 0x1000 725 ))?, 726 self.snapshot().revision(), 727 input, 728 DEFAULT_COMMAND_TIMEOUT, 729 ) 730 .await 731 } 732 733 #[cfg(test)] 734 async fn acknowledge_generated_key_stage_test( 735 &self, 736 id: RecoveryStageId, 737 ) -> Result<AppSnapshot, SafeError> { 738 let request_number = self.next_request.fetch_add(1, Ordering::Relaxed); 739 self.acknowledge_generated_key_stage( 740 id, 741 DurableRequestId::parse(format!( 742 "01890f3e-7b1c-7000-8000-{:012x}", 743 request_number + 0x2000 744 ))?, 745 self.snapshot().revision(), 746 DEFAULT_COMMAND_TIMEOUT, 747 ) 748 .await 749 } 750 751 #[cfg(test)] 752 async fn confirm_identity_removal_test( 753 &self, 754 token: RemovalConfirmationToken, 755 ) -> Result<AppSnapshot, SafeError> { 756 let request_number = self.next_request.fetch_add(1, Ordering::Relaxed); 757 self.confirm_identity_removal( 758 token, 759 DurableRequestId::parse(format!( 760 "01890f3e-7b1c-7000-8000-{:012x}", 761 request_number + 0x3000 762 ))?, 763 self.snapshot().revision(), 764 DEFAULT_COMMAND_TIMEOUT, 765 ) 766 .await 767 } 768 769 #[cfg(test)] 770 async fn import_secret_key_with_timeout( 771 &self, 772 input: SecretKeyInput, 773 timeout: Duration, 774 ) -> Result<ImportIdentityReceipt, SafeError> { 775 let raw_request = self.next_request.fetch_add(1, Ordering::Relaxed); 776 let request_id = RequestId::new(raw_request).ok_or_else(request_space_exhausted)?; 777 let expected_revision = self.adapter.core().snapshot().revision(); 778 let durable_request = DurableRequestId::parse(format!( 779 "01890f3e-7b1c-7000-8000-{:012x}", 780 raw_request + 0x4000 781 ))?; 782 match self 783 .dispatch_with_deadline( 784 RuntimeCommand::ImportSecretKey { 785 input, 786 durable_request, 787 expected_revision: expected_revision.value(), 788 }, 789 Some(expected_revision), 790 request_id, 791 Instant::now() + timeout, 792 ) 793 .await? 794 { 795 RuntimeCommandValue::Imported(receipt) => Ok(receipt), 796 _ => Err(invalid_actor_response()), 797 } 798 } 799 800 fn expect_snapshot(value: RuntimeCommandValue) -> Result<AppSnapshot, SafeError> { 801 match value { 802 RuntimeCommandValue::Snapshot(snapshot) => Ok(*snapshot), 803 _ => Err(invalid_actor_response()), 804 } 805 } 806 } 807 808 impl RuntimeActor { 809 async fn run( 810 mut self, 811 mut receiver: mpsc::Receiver<CommandEnvelope<RuntimeCommand, RuntimeCommandValue>>, 812 ) { 813 let (completion_sender, mut completions) = mpsc::channel(DEFAULT_TASK_CAPACITY); 814 loop { 815 tokio::select! { 816 envelope = receiver.recv() => { 817 let Some(envelope) = envelope else { 818 break; 819 }; 820 if !self.handle_command(envelope, &completion_sender).await { 821 break; 822 } 823 } 824 completion = completions.recv(), if !self.profile_tasks.is_empty() => { 825 if let Some(completion) = completion { 826 self.complete_profile_task(completion).await; 827 } 828 } 829 } 830 } 831 self.cancel_profile_tasks(None).await; 832 } 833 834 async fn handle_command( 835 &mut self, 836 envelope: CommandEnvelope<RuntimeCommand, RuntimeCommandValue>, 837 completion_sender: &mpsc::Sender<ProfileCompletion>, 838 ) -> bool { 839 let (context, command, reply) = envelope.into_parts(); 840 if matches!(command, RuntimeCommand::SubscribeChanges(_)) && reply.is_closed() { 841 return true; 842 } 843 if let Some(result) = self.preflight(context, &command) { 844 let _ = reply.send(CommandReceipt::new(context.request_id(), result)); 845 return true; 846 } 847 if matches!(command, RuntimeCommand::RefreshActiveProfile) { 848 self.start_profile_task(context, reply, completion_sender.clone()) 849 .await; 850 return true; 851 } 852 if matches!(command, RuntimeCommand::Close) { 853 let result = self.close_actor().await; 854 let closed = matches!(result, CommandResult::Completed(_)); 855 let _ = reply.send(CommandReceipt::new(context.request_id(), result)); 856 return !closed; 857 } 858 let changes_session = matches!( 859 command, 860 RuntimeCommand::ActivateIdentity(_) 861 | RuntimeCommand::SignOut 862 | RuntimeCommand::ConfirmIdentityRemoval { .. } 863 ); 864 let begins_generated_recovery = matches!(&command, RuntimeCommand::BeginGeneratedKeyStage); 865 let result = self.execute_command(context, command).await; 866 if begins_generated_recovery && matches!(&result, CommandResult::Completed(_)) { 867 let snapshot = self.adapter.core().snapshot(); 868 self.cancel_profile_tasks(Some(&snapshot)).await; 869 } 870 if changes_session && matches!(result, CommandResult::Completed(_)) { 871 self.advance_session_generation().await; 872 self.synchronize_foreground_session(); 873 } 874 if matches!(result, CommandResult::Completed(_)) { 875 self.changes.publish(self.adapter.core().snapshot()); 876 } 877 if let Err(receipt) = reply.send(CommandReceipt::new(context.request_id(), result)) 878 && let CommandResult::Completed(RuntimeCommandValue::Subscription(subscription)) = 879 receipt.into_result() 880 { 881 let _ = self.changes.unsubscribe(subscription.id()); 882 } 883 true 884 } 885 886 fn preflight( 887 &self, 888 context: CommandContext, 889 command: &RuntimeCommand, 890 ) -> Option<CommandResult<RuntimeCommandValue>> { 891 if context.is_expired(Instant::now()) { 892 return Some(CommandResult::TimedOut); 893 } 894 let lifecycle = match self.lifecycle.lock() { 895 Ok(lifecycle) => lifecycle.to_owned(), 896 Err(_) => return Some(CommandResult::Failed(runtime_state_unavailable())), 897 }; 898 if matches!(lifecycle.lifecycle(), RuntimeLifecycle::Closed) { 899 return Some(CommandResult::Closed); 900 } 901 if !lifecycle.allows(command.class()) { 902 return Some(CommandResult::Failed(command_unavailable())); 903 } 904 if self.generated_key_stage.pending().is_some() 905 && !matches!( 906 command, 907 RuntimeCommand::Snapshot 908 | RuntimeCommand::AcknowledgeGeneratedKeyStage { .. } 909 | RuntimeCommand::CancelGeneratedKeyStage 910 | RuntimeCommand::SubscribeChanges(_) 911 | RuntimeCommand::UnsubscribeChanges(_) 912 | RuntimeCommand::Close 913 ) 914 { 915 return Some(CommandResult::Failed(generated_recovery_route_active())); 916 } 917 let current_revision = self.adapter.core().snapshot().revision(); 918 if !command.resolves_revision_through_durable_replay() 919 && context 920 .expected_revision() 921 .is_some_and(|expected| expected != current_revision) 922 { 923 return Some(CommandResult::Conflicted { current_revision }); 924 } 925 None 926 } 927 928 async fn execute_command( 929 &mut self, 930 context: CommandContext, 931 command: RuntimeCommand, 932 ) -> CommandResult<RuntimeCommandValue> { 933 let result = match command { 934 RuntimeCommand::Snapshot => Ok(RuntimeCommandValue::Snapshot(Box::new( 935 self.adapter.core().snapshot(), 936 ))), 937 RuntimeCommand::GenerateIdentity { 938 durable_request, 939 expected_revision, 940 } => { 941 self.run_async( 942 context.deadline(), 943 move |adapter, secrets, clock| async move { 944 adapter 945 .generate_identity_durable( 946 &durable_request, 947 expected_revision, 948 secrets.as_ref(), 949 clock.as_ref(), 950 ) 951 .await 952 .map(RuntimeCommandValue::Generated) 953 }, 954 ) 955 .await 956 } 957 RuntimeCommand::BeginGeneratedKeyStage => { 958 match NonZeroU64::new(context.request_id().get()).map(RecoveryStageId::new) { 959 Some(stage_id) => self 960 .generated_key_stage 961 .begin( 962 self.adapter.key_material(), 963 stage_id, 964 self.adapter.core().snapshot().revision().value(), 965 self.clock.now(), 966 ) 967 .map(RuntimeCommandValue::GeneratedKeyStage), 968 None => Err(request_space_exhausted()), 969 } 970 } 971 RuntimeCommand::AcknowledgeGeneratedKeyStage { 972 id, 973 durable_request, 974 } => match self.generated_key_stage.take(id, self.clock.now()) { 975 Ok(staged) => { 976 self.run_async( 977 context.deadline(), 978 move |adapter, secrets, clock| async move { 979 commit_generated_key_stage( 980 adapter.as_ref(), 981 secrets.as_ref(), 982 clock.as_ref(), 983 &durable_request, 984 staged, 985 ) 986 .await 987 }, 988 ) 989 .await 990 } 991 Err(error) => Err(error), 992 }, 993 RuntimeCommand::CancelGeneratedKeyStage => Ok( 994 RuntimeCommandValue::GeneratedKeyStageCancelled(self.generated_key_stage.cancel()), 995 ), 996 RuntimeCommand::ImportSecretKey { 997 input, 998 durable_request, 999 expected_revision, 1000 } => { 1001 self.run_async( 1002 context.deadline(), 1003 move |adapter, secrets, clock| async move { 1004 adapter 1005 .import_secret_key_durable( 1006 &durable_request, 1007 expected_revision, 1008 input, 1009 secrets.as_ref(), 1010 clock.as_ref(), 1011 ) 1012 .await 1013 .map(RuntimeCommandValue::Imported) 1014 }, 1015 ) 1016 .await 1017 } 1018 RuntimeCommand::SelectIdentity(public_key) => { 1019 self.run_async(context.deadline(), move |adapter, _, _| async move { 1020 adapter 1021 .select_identity(public_key) 1022 .await 1023 .map(Box::new) 1024 .map(RuntimeCommandValue::Snapshot) 1025 }) 1026 .await 1027 } 1028 RuntimeCommand::ActivateIdentity(public_key) => { 1029 self.run_async( 1030 context.deadline(), 1031 move |adapter, secrets, clock| async move { 1032 adapter 1033 .activate_identity(public_key, secrets.as_ref(), clock.as_ref()) 1034 .await 1035 .map(Box::new) 1036 .map(RuntimeCommandValue::Snapshot) 1037 }, 1038 ) 1039 .await 1040 } 1041 RuntimeCommand::SignOut => self 1042 .adapter 1043 .sign_out() 1044 .map(Box::new) 1045 .map(RuntimeCommandValue::Snapshot), 1046 RuntimeCommand::RequestIdentityRemoval(public_key) => self 1047 .adapter 1048 .request_identity_removal(public_key, self.clock.as_ref()) 1049 .map(RuntimeCommandValue::RemovalRequest), 1050 RuntimeCommand::ConfirmIdentityRemoval { 1051 token, 1052 durable_request, 1053 } => { 1054 self.run_async( 1055 context.deadline(), 1056 move |adapter, secrets, clock| async move { 1057 adapter 1058 .confirm_identity_removal_durable( 1059 &durable_request, 1060 token, 1061 secrets.as_ref(), 1062 clock.as_ref(), 1063 ) 1064 .await 1065 .map(Box::new) 1066 .map(RuntimeCommandValue::Snapshot) 1067 }, 1068 ) 1069 .await 1070 } 1071 RuntimeCommand::SubscribeChanges(capacity) => self 1072 .changes 1073 .subscribe(capacity) 1074 .map(|(id, receiver)| { 1075 RuntimeCommandValue::Subscription(RuntimeChangeSubscription { id, receiver }) 1076 }) 1077 .ok_or_else(observer_registration_failed), 1078 RuntimeCommand::UnsubscribeChanges(id) => Ok(RuntimeCommandValue::Unsubscribed( 1079 self.changes.unsubscribe(id), 1080 )), 1081 RuntimeCommand::Close | RuntimeCommand::RefreshActiveProfile => { 1082 Err(invalid_actor_response()) 1083 } 1084 }; 1085 result.map_or_else(CommandResult::Failed, CommandResult::Completed) 1086 } 1087 1088 async fn run_async<F, Fut>( 1089 &self, 1090 deadline: Instant, 1091 operation: F, 1092 ) -> Result<RuntimeCommandValue, SafeError> 1093 where 1094 F: FnOnce(Arc<PersistentAppCore>, Arc<dyn SecretStore>, Arc<dyn Clock>) -> Fut, 1095 Fut: Future<Output = Result<RuntimeCommandValue, SafeError>>, 1096 { 1097 let adapter = Arc::clone(&self.adapter); 1098 let secrets = Arc::clone(&self.secrets); 1099 let clock = Arc::clone(&self.clock); 1100 let remaining = deadline.saturating_duration_since(Instant::now()); 1101 tokio::time::timeout(remaining, operation(adapter, secrets, clock)) 1102 .await 1103 .map_err(|_| command_timed_out())? 1104 } 1105 1106 async fn start_profile_task( 1107 &mut self, 1108 context: CommandContext, 1109 reply: oneshot::Sender<CommandReceipt<RuntimeCommandValue>>, 1110 completion_sender: mpsc::Sender<ProfileCompletion>, 1111 ) { 1112 let foreground = match self.published_foreground_session.lock() { 1113 Ok(foreground) => foreground.clone(), 1114 Err(_) => { 1115 let _ = reply.send(CommandReceipt::new( 1116 context.request_id(), 1117 CommandResult::Failed(runtime_state_unavailable()), 1118 )); 1119 return; 1120 } 1121 }; 1122 let plan = match self.adapter.core().begin_profile_refresh() { 1123 Ok(Some(plan)) => plan, 1124 Ok(None) => { 1125 let _ = reply.send(CommandReceipt::new( 1126 context.request_id(), 1127 CommandResult::Completed(RuntimeCommandValue::Snapshot(Box::new( 1128 self.adapter.core().snapshot(), 1129 ))), 1130 )); 1131 return; 1132 } 1133 Err(error) => { 1134 let _ = reply.send(CommandReceipt::new( 1135 context.request_id(), 1136 CommandResult::Failed(error), 1137 )); 1138 return; 1139 } 1140 }; 1141 let Some(foreground) = foreground.filter(|binding| { 1142 binding.identity().public_key() == plan.public_key() 1143 && binding.generation() == self.session_generation 1144 }) else { 1145 let _ = reply.send(CommandReceipt::new( 1146 context.request_id(), 1147 CommandResult::Failed(stale_profile_binding()), 1148 )); 1149 return; 1150 }; 1151 let correlation = TaskCorrelation::new( 1152 context.request_id(), 1153 plan.public_key(), 1154 foreground.signer_binding(), 1155 plan.expected_revision(), 1156 self.session_generation, 1157 ); 1158 let client = Arc::clone(&self.nostr); 1159 let relays = plan.relays().to_vec(); 1160 let request_id = context.request_id(); 1161 let handle = self.runtime.spawn(async move { 1162 let result = client 1163 .fetch_profile(correlation.identity(), &relays, context.deadline()) 1164 .await; 1165 let _ = completion_sender 1166 .send(ProfileCompletion { request_id, result }) 1167 .await; 1168 }); 1169 let previous = self.profile_tasks.insert( 1170 request_id, 1171 PendingProfileTask { 1172 correlation, 1173 plan, 1174 deadline: context.deadline(), 1175 reply, 1176 handle, 1177 }, 1178 ); 1179 if let Some(previous) = previous { 1180 self.fail_lifecycle(request_space_exhausted()); 1181 previous.handle.abort(); 1182 let _ = previous.handle.await; 1183 let _ = previous.reply.send(CommandReceipt::new( 1184 previous.correlation.request_id(), 1185 CommandResult::Failed(request_space_exhausted()), 1186 )); 1187 self.cancel_profile_tasks(None).await; 1188 } 1189 } 1190 1191 async fn close_actor(&mut self) -> CommandResult<RuntimeCommandValue> { 1192 let begin = { 1193 let Ok(mut lifecycle) = self.lifecycle.lock() else { 1194 return CommandResult::Failed(runtime_state_unavailable()); 1195 }; 1196 lifecycle.begin_shutdown() 1197 }; 1198 match begin { 1199 Ok(()) => { 1200 self.generated_key_stage.cancel(); 1201 self.cancel_profile_tasks(None).await; 1202 if let Err(error) = self.adapter.close().await { 1203 self.fail_lifecycle(error); 1204 return CommandResult::Failed(error); 1205 } 1206 self.changes.close(); 1207 let Ok(mut foreground) = self.published_foreground_session.lock() else { 1208 return CommandResult::Failed(runtime_state_unavailable()); 1209 }; 1210 *foreground = None; 1211 drop(foreground); 1212 let Ok(mut lifecycle) = self.lifecycle.lock() else { 1213 return CommandResult::Failed(runtime_state_unavailable()); 1214 }; 1215 match lifecycle.finish_shutdown() { 1216 Ok(()) => CommandResult::Completed(RuntimeCommandValue::Closed), 1217 Err(error) => CommandResult::Failed(error), 1218 } 1219 } 1220 Err(error) => CommandResult::Failed(error), 1221 } 1222 } 1223 1224 async fn complete_profile_task(&mut self, completion: ProfileCompletion) { 1225 let Some(task) = self.profile_tasks.remove(&completion.request_id) else { 1226 return; 1227 }; 1228 let _ = task.handle.await; 1229 let current = self.adapter.core().snapshot(); 1230 let foreground = match self.published_foreground_session.lock() { 1231 Ok(foreground) => foreground.clone(), 1232 Err(_) => { 1233 self.fail_lifecycle(runtime_state_unavailable()); 1234 let _ = task.reply.send(CommandReceipt::new( 1235 task.correlation.request_id(), 1236 CommandResult::Failed(runtime_state_unavailable()), 1237 )); 1238 return; 1239 } 1240 }; 1241 let correlated = task.correlation.session_generation() == self.session_generation 1242 && foreground.is_some_and(|binding| { 1243 binding.generation() == task.correlation.session_generation() 1244 && binding.identity().public_key() == task.correlation.identity() 1245 && binding.signer_binding() == task.correlation.binding() 1246 }) 1247 && current.active_identity().is_some_and(|active| { 1248 active.identity().public_key() == task.correlation.identity() 1249 }); 1250 let result = if correlated { 1251 let plan = task.plan.clone(); 1252 let completed = self 1253 .run_async(task.deadline, move |adapter, _, clock| async move { 1254 adapter 1255 .core() 1256 .complete_profile_refresh( 1257 &plan, 1258 completion.result, 1259 adapter.database(), 1260 clock.as_ref(), 1261 ) 1262 .await 1263 .map(Box::new) 1264 .map(RuntimeCommandValue::Snapshot) 1265 }) 1266 .await; 1267 completed.map_or_else(CommandResult::Failed, CommandResult::Completed) 1268 } else { 1269 CommandResult::Completed(RuntimeCommandValue::Snapshot(Box::new(current))) 1270 }; 1271 if matches!(result, CommandResult::Completed(_)) { 1272 self.changes.publish(self.adapter.core().snapshot()); 1273 } 1274 let _ = task 1275 .reply 1276 .send(CommandReceipt::new(task.correlation.request_id(), result)); 1277 } 1278 1279 async fn advance_session_generation(&mut self) { 1280 let Some(next) = self.session_generation.next() else { 1281 self.fail_lifecycle(request_space_exhausted()); 1282 self.cancel_profile_tasks(None).await; 1283 return; 1284 }; 1285 self.session_generation = next; 1286 self.published_session_generation 1287 .store(next.value(), Ordering::Release); 1288 let snapshot = self.adapter.core().snapshot(); 1289 self.cancel_profile_tasks(Some(&snapshot)).await; 1290 } 1291 1292 fn synchronize_foreground_session(&mut self) { 1293 let session = self 1294 .adapter 1295 .core() 1296 .snapshot() 1297 .active_identity() 1298 .map(|active| { 1299 let public_key = active.identity().public_key(); 1300 ActiveSessionBinding::new( 1301 NostrIdentityReference::derive(public_key)?, 1302 LocalKeyringBinding::new(public_key, SignerAvailability::Available), 1303 self.session_generation, 1304 ) 1305 }); 1306 let session = match session.transpose() { 1307 Ok(session) => session, 1308 Err(error) => { 1309 self.fail_lifecycle(error); 1310 None 1311 } 1312 }; 1313 match self.published_foreground_session.lock() { 1314 Ok(mut foreground) => *foreground = session, 1315 Err(_) => self.fail_lifecycle(runtime_state_unavailable()), 1316 } 1317 } 1318 1319 fn fail_lifecycle(&self, error: SafeError) { 1320 if let Ok(mut lifecycle) = self.lifecycle.lock() { 1321 lifecycle.fail(error); 1322 } 1323 } 1324 1325 async fn cancel_profile_tasks(&mut self, snapshot: Option<&AppSnapshot>) { 1326 let tasks = std::mem::take(&mut self.profile_tasks); 1327 for (_, task) in tasks { 1328 task.handle.abort(); 1329 let _ = task.handle.await; 1330 let receipt_result = snapshot.map_or(CommandResult::Closed, |snapshot| { 1331 CommandResult::Completed(RuntimeCommandValue::Snapshot(Box::new(snapshot.clone()))) 1332 }); 1333 let _ = task.reply.send(CommandReceipt::new( 1334 task.correlation.request_id(), 1335 receipt_result, 1336 )); 1337 } 1338 } 1339 } 1340 1341 async fn commit_generated_key_stage( 1342 adapter: &PersistentAppCore, 1343 secrets: &dyn SecretStore, 1344 clock: &dyn Clock, 1345 request: &DurableRequestId, 1346 staged: StagedGeneratedKey, 1347 ) -> Result<RuntimeCommandValue, SafeError> { 1348 adapter 1349 .commit_staged_generated_key(request, staged, secrets, clock) 1350 .await?; 1351 Ok(RuntimeCommandValue::Snapshot(Box::new( 1352 adapter.core().snapshot(), 1353 ))) 1354 } 1355 1356 #[cfg(test)] 1357 fn test_runtime_context(root: &std::path::Path) -> Result<RuntimeContext, SafeError> { 1358 let canonical = root 1359 .canonicalize() 1360 .map_err(|_| invalid_runtime_evidence())?; 1361 let resolver = RadrootsPathResolver::new( 1362 RadrootsPlatform::current(), 1363 RadrootsHostEnvironment::default(), 1364 ); 1365 let bootstrap = RuntimeContextBootstrap::new( 1366 RadrootsPathProfile::RepoLocal, 1367 Some(canonical), 1368 RuntimeContextSource::BootstrapCli, 1369 RuntimeContextSource::SafeDefault, 1370 ) 1371 .map_err(|_| invalid_runtime_evidence())?; 1372 let context = RuntimeContext::resolve( 1373 &resolver, 1374 bootstrap, 1375 ServiceId::new("harvestcircle").map_err(|_| invalid_runtime_evidence())?, 1376 InstanceId::new("desktop").map_err(|_| invalid_runtime_evidence())?, 1377 ) 1378 .map_err(|_| invalid_runtime_evidence())?; 1379 std::fs::create_dir_all(root.join("data")).map_err(|_| invalid_runtime_evidence())?; 1380 Ok(context) 1381 } 1382 1383 #[cfg(test)] 1384 fn test_migration_build_identity() -> Result<MigrationBuildIdentity, SafeError> { 1385 MigrationBuildIdentity::new( 1386 "0.1.0-alpha", 1387 "1111111111111111111111111111111111111111", 1388 "2222222222222222222222222222222222222222", 1389 "1.97.1", 1390 "test", 1391 "test", 1392 1, 1393 1, 1394 1, 1395 1, 1396 1, 1397 ) 1398 .map_err(|_| invalid_runtime_evidence()) 1399 } 1400 1401 const fn request_space_exhausted() -> SafeError { 1402 SafeError::new( 1403 SafeErrorCode::InvalidApplicationState, 1404 SafeMessage::new("The runtime request identifier space is exhausted."), 1405 ) 1406 } 1407 1408 const fn invalid_runtime_evidence() -> SafeError { 1409 SafeError::new( 1410 SafeErrorCode::InvalidApplicationState, 1411 SafeMessage::new("The runtime startup evidence is invalid."), 1412 ) 1413 } 1414 1415 const fn stale_profile_binding() -> SafeError { 1416 SafeError::new( 1417 SafeErrorCode::InvalidApplicationState, 1418 SafeMessage::new("The active identity binding changed before profile refresh."), 1419 ) 1420 } 1421 1422 const fn command_conflicted() -> SafeError { 1423 SafeError::new( 1424 SafeErrorCode::InvalidApplicationState, 1425 SafeMessage::new("The command conflicts with newer application state."), 1426 ) 1427 } 1428 1429 const fn command_rejected() -> SafeError { 1430 SafeError::new( 1431 SafeErrorCode::InvalidApplicationState, 1432 SafeMessage::new("The runtime is busy. Try again."), 1433 ) 1434 } 1435 1436 const fn command_timed_out() -> SafeError { 1437 SafeError::new( 1438 SafeErrorCode::InvalidApplicationState, 1439 SafeMessage::new("The runtime command timed out."), 1440 ) 1441 } 1442 1443 const fn runtime_closed() -> SafeError { 1444 SafeError::new( 1445 SafeErrorCode::InvalidApplicationState, 1446 SafeMessage::new("The application runtime is closed."), 1447 ) 1448 } 1449 1450 const fn runtime_state_unavailable() -> SafeError { 1451 SafeError::new( 1452 SafeErrorCode::InvalidApplicationState, 1453 SafeMessage::new("The application runtime state is unavailable."), 1454 ) 1455 } 1456 1457 const fn command_unavailable() -> SafeError { 1458 SafeError::new( 1459 SafeErrorCode::InvalidApplicationState, 1460 SafeMessage::new("The command is unavailable in the current runtime state."), 1461 ) 1462 } 1463 1464 const fn generated_recovery_route_active() -> SafeError { 1465 SafeError::new( 1466 SafeErrorCode::InvalidApplicationState, 1467 SafeMessage::new("Complete or cancel generated-key recovery before another action."), 1468 ) 1469 } 1470 1471 const fn invalid_actor_response() -> SafeError { 1472 SafeError::new( 1473 SafeErrorCode::InvalidApplicationState, 1474 SafeMessage::new("The runtime returned an invalid command response."), 1475 ) 1476 } 1477 1478 const fn observer_registration_failed() -> SafeError { 1479 SafeError::new( 1480 SafeErrorCode::ObserverRegistrationFailed, 1481 SafeMessage::new("The application change subscription could not be registered."), 1482 ) 1483 } 1484 1485 #[cfg(test)] 1486 mod tests { 1487 use std::future::Future; 1488 use std::num::NonZeroUsize; 1489 use std::sync::Arc; 1490 use std::sync::atomic::{AtomicBool, Ordering}; 1491 use std::task::{Context, Poll, Wake, Waker}; 1492 use std::thread::{self, Thread}; 1493 use std::time::{Duration, Instant}; 1494 1495 use harvestcircle_application::{ 1496 ActiveSessionBinding, ActorMailbox, BoxFuture, Clock, CommandContext, CommandSubmission, 1497 DurableRequestId, FailureSecretStore, GeneratedKeyStage, InMemorySecretStore, NostrClient, 1498 OrderedSnapshotChanges, ProfileFetchResult, RelayConfiguration, RequestId, 1499 RuntimeLifecycle, SecretStore, SecretStoreOperation, SessionGeneration, SessionState, 1500 SnapshotRevision, 1501 }; 1502 use harvestcircle_domain::{ 1503 LocalKeyringBinding, NostrIdentityReference, PublicKey, SafeError, SafeErrorCode, 1504 SecretKeyInput, SignerAvailability, UnixTimestamp, 1505 }; 1506 use radroots_transport_nostr::{RelayAccess, RelayEndpoint, RelayUrlPolicy}; 1507 1508 use super::{ 1509 DEFAULT_COMMAND_TIMEOUT, RuntimeActor, RuntimeActorHandle, RuntimeCommand, 1510 RuntimeDependencies, command_unavailable, test_migration_build_identity, 1511 test_runtime_context, 1512 }; 1513 use crate::{InstallationIdentity, InstallationIdentitySource, UuidInstallationIdentitySource}; 1514 1515 struct FixedInstallationIdentity(&'static str); 1516 1517 impl InstallationIdentitySource for FixedInstallationIdentity { 1518 fn generate(&self) -> Result<InstallationIdentity, SafeError> { 1519 InstallationIdentity::parse(self.0) 1520 } 1521 } 1522 1523 fn dependencies( 1524 secrets: Arc<dyn SecretStore>, 1525 nostr: Arc<dyn NostrClient>, 1526 ) -> RuntimeDependencies { 1527 RuntimeDependencies::new( 1528 secrets, 1529 Arc::new(FixedClock), 1530 nostr, 1531 Arc::new(UuidInstallationIdentitySource), 1532 ) 1533 } 1534 1535 struct FixedClock; 1536 1537 impl Clock for FixedClock { 1538 fn now(&self) -> UnixTimestamp { 1539 UnixTimestamp::from_seconds(50).expect("time") 1540 } 1541 } 1542 1543 struct OfflineNostr; 1544 1545 impl NostrClient for OfflineNostr { 1546 fn fetch_profile<'a>( 1547 &'a self, 1548 _public_key: PublicKey, 1549 _relays: &'a [RelayEndpoint], 1550 _deadline: Instant, 1551 ) -> BoxFuture<'a, Result<ProfileFetchResult, SafeError>> { 1552 Box::pin(async { Ok(ProfileFetchResult::complete(None)) }) 1553 } 1554 } 1555 1556 struct BlockingNostr { 1557 started: tokio::sync::Semaphore, 1558 release: tokio::sync::Semaphore, 1559 } 1560 1561 impl BlockingNostr { 1562 fn new() -> Self { 1563 Self { 1564 started: tokio::sync::Semaphore::new(0), 1565 release: tokio::sync::Semaphore::new(0), 1566 } 1567 } 1568 } 1569 1570 impl NostrClient for BlockingNostr { 1571 fn fetch_profile<'a>( 1572 &'a self, 1573 _public_key: PublicKey, 1574 _relays: &'a [RelayEndpoint], 1575 _deadline: Instant, 1576 ) -> BoxFuture<'a, Result<ProfileFetchResult, SafeError>> { 1577 Box::pin(async move { 1578 self.started.add_permits(1); 1579 let permit = self.release.acquire().await.expect("release"); 1580 permit.forget(); 1581 Ok(ProfileFetchResult::complete(None)) 1582 }) 1583 } 1584 } 1585 1586 struct BlockingSecretStore { 1587 inner: InMemorySecretStore, 1588 block_next_put: AtomicBool, 1589 put_started: AtomicBool, 1590 released: AtomicBool, 1591 release_signal: tokio::sync::Notify, 1592 } 1593 1594 impl BlockingSecretStore { 1595 fn new() -> Self { 1596 Self { 1597 inner: InMemorySecretStore::default(), 1598 block_next_put: AtomicBool::new(true), 1599 put_started: AtomicBool::new(false), 1600 released: AtomicBool::new(false), 1601 release_signal: tokio::sync::Notify::new(), 1602 } 1603 } 1604 1605 async fn wait_until_put_started(&self) { 1606 while !self.put_started.load(Ordering::Acquire) { 1607 tokio::task::yield_now().await; 1608 } 1609 } 1610 1611 fn release(&self) { 1612 self.released.store(true, Ordering::Release); 1613 self.release_signal.notify_waiters(); 1614 } 1615 } 1616 1617 impl SecretStore for BlockingSecretStore { 1618 fn put<'a>( 1619 &'a self, 1620 request_id: &'a harvestcircle_application::DurableRequestId, 1621 public_key: PublicKey, 1622 secret: SecretKeyInput, 1623 ) -> BoxFuture<'a, Result<(), SafeError>> { 1624 Box::pin(async move { 1625 if self.block_next_put.swap(false, Ordering::AcqRel) { 1626 self.put_started.store(true, Ordering::Release); 1627 loop { 1628 let notified = self.release_signal.notified(); 1629 if self.released.load(Ordering::Acquire) { 1630 break; 1631 } 1632 notified.await; 1633 } 1634 } 1635 self.inner.put(request_id, public_key, secret).await 1636 }) 1637 } 1638 1639 fn load(&self, public_key: PublicKey) -> BoxFuture<'_, Result<SecretKeyInput, SafeError>> { 1640 self.inner.load(public_key) 1641 } 1642 1643 fn contains(&self, public_key: PublicKey) -> BoxFuture<'_, Result<bool, SafeError>> { 1644 self.inner.contains(public_key) 1645 } 1646 1647 fn delete<'a>( 1648 &'a self, 1649 request_id: &'a harvestcircle_application::DurableRequestId, 1650 public_key: PublicKey, 1651 ) -> BoxFuture<'a, Result<(), SafeError>> { 1652 self.inner.delete(request_id, public_key) 1653 } 1654 } 1655 1656 async fn actor() -> (RuntimeActorHandle, Arc<InMemorySecretStore>) { 1657 let secrets = Arc::new(InMemorySecretStore::default()); 1658 let secret_port: Arc<dyn SecretStore> = secrets.clone(); 1659 let actor = RuntimeActorHandle::in_memory( 1660 RelayConfiguration::default(), 1661 dependencies(secret_port, Arc::new(OfflineNostr)), 1662 NonZeroUsize::new(8).expect("capacity"), 1663 &tokio::runtime::Handle::current(), 1664 ) 1665 .await 1666 .expect("actor"); 1667 (actor, secrets) 1668 } 1669 1670 struct ThreadWake(Thread); 1671 1672 impl Wake for ThreadWake { 1673 fn wake(self: Arc<Self>) { 1674 self.0.unpark(); 1675 } 1676 1677 fn wake_by_ref(self: &Arc<Self>) { 1678 self.0.unpark(); 1679 } 1680 } 1681 1682 fn block_on_without_runtime<F: Future>(future: F) -> F::Output { 1683 let waker = Waker::from(Arc::new(ThreadWake(thread::current()))); 1684 let mut context = Context::from_waker(&waker); 1685 let mut future = std::pin::pin!(future); 1686 loop { 1687 match future.as_mut().poll(&mut context) { 1688 Poll::Ready(output) => return output, 1689 Poll::Pending => thread::park(), 1690 } 1691 } 1692 } 1693 1694 #[test] 1695 fn actor_operations_support_foreign_executor_polling() { 1696 let runtime = tokio::runtime::Runtime::new().expect("runtime"); 1697 let (actor, _) = runtime.block_on(actor()); 1698 1699 let snapshot = block_on_without_runtime(actor.bootstrap()).expect("bootstrap"); 1700 assert_eq!(snapshot, actor.snapshot()); 1701 block_on_without_runtime(actor.close()).expect("close"); 1702 } 1703 1704 #[tokio::test(flavor = "multi_thread")] 1705 async fn installation_identity_survives_file_backed_runtime_restart() { 1706 let directory = tempfile::tempdir().expect("temporary directory"); 1707 let context = test_runtime_context(directory.path()).expect("runtime context"); 1708 let build = test_migration_build_identity().expect("build identity"); 1709 let first = RuntimeActorHandle::open( 1710 &context, 1711 RelayConfiguration::default(), 1712 RuntimeDependencies::new( 1713 Arc::new(InMemorySecretStore::default()), 1714 Arc::new(FixedClock), 1715 Arc::new(OfflineNostr), 1716 Arc::new(FixedInstallationIdentity( 1717 "11aabbccddeeff001122334455667788", 1718 )), 1719 ), 1720 &build, 1721 NonZeroUsize::new(8).expect("capacity"), 1722 &tokio::runtime::Handle::current(), 1723 ) 1724 .await 1725 .expect("first runtime"); 1726 assert_eq!( 1727 first.installation_identity().as_str(), 1728 "11aabbccddeeff001122334455667788" 1729 ); 1730 first.close().await.expect("first close"); 1731 drop(first); 1732 1733 let second = RuntimeActorHandle::open( 1734 &context, 1735 RelayConfiguration::default(), 1736 RuntimeDependencies::new( 1737 Arc::new(InMemorySecretStore::default()), 1738 Arc::new(FixedClock), 1739 Arc::new(OfflineNostr), 1740 Arc::new(FixedInstallationIdentity( 1741 "22aabbccddeeff001122334455667788", 1742 )), 1743 ), 1744 &build, 1745 NonZeroUsize::new(8).expect("capacity"), 1746 &tokio::runtime::Handle::current(), 1747 ) 1748 .await 1749 .expect("second runtime"); 1750 assert_eq!( 1751 second.installation_identity().as_str(), 1752 "11aabbccddeeff001122334455667788" 1753 ); 1754 second.close().await.expect("second close"); 1755 } 1756 1757 #[tokio::test(flavor = "multi_thread")] 1758 async fn identity_mutations_run_serially_through_one_ready_actor() { 1759 let (actor, secrets) = actor().await; 1760 assert_eq!(actor.lifecycle(), RuntimeLifecycle::Ready); 1761 1762 let imported = actor 1763 .import_secret_key_test( 1764 SecretKeyInput::parse( 1765 "7e7e9c42a91bfef19fa7ea99d52d8afdb67d893a8fefba1f5cb9793f2107f6d7".to_owned(), 1766 ) 1767 .expect("input"), 1768 ) 1769 .await 1770 .expect("import"); 1771 let public_key = imported.identity().public_key(); 1772 let activated = actor.activate_identity(public_key).await.expect("activate"); 1773 assert_eq!(activated.session(), SessionState::Active); 1774 let foreground = actor.foreground_session().expect("foreground session"); 1775 assert_eq!(foreground.identity().public_key(), public_key); 1776 assert_eq!(foreground.signer_binding().identity(), public_key); 1777 assert_eq!(foreground.generation(), actor.session_generation()); 1778 assert!(secrets.contains(public_key).await.expect("credential")); 1779 1780 let signed_out = actor.sign_out().await.expect("sign out"); 1781 assert_eq!(signed_out.session(), SessionState::SignedOut); 1782 assert!(actor.foreground_session().is_none()); 1783 let removal = actor 1784 .request_identity_removal(public_key) 1785 .await 1786 .expect("removal request"); 1787 let removed = actor 1788 .confirm_identity_removal_test(removal) 1789 .await 1790 .expect("remove"); 1791 assert!(removed.identities().is_empty()); 1792 assert!( 1793 !secrets 1794 .contains(public_key) 1795 .await 1796 .expect("credential removed") 1797 ); 1798 } 1799 1800 #[tokio::test(flavor = "multi_thread")] 1801 async fn public_actor_commands_cover_generation_selection_and_empty_profile_refresh() { 1802 let (actor, _) = actor().await; 1803 let unchanged = actor 1804 .refresh_active_profile() 1805 .await 1806 .expect("refresh without an active identity"); 1807 assert!(unchanged.active_identity().is_none()); 1808 1809 let generated = actor 1810 .generate_identity( 1811 DurableRequestId::parse("01890f3e-7b1c-7000-8000-000000005001").expect("request"), 1812 actor.snapshot().revision(), 1813 DEFAULT_COMMAND_TIMEOUT, 1814 ) 1815 .await 1816 .expect("generate identity"); 1817 let selected = actor 1818 .select_identity(generated.identity().public_key()) 1819 .await 1820 .expect("select generated identity"); 1821 assert_eq!( 1822 selected.selected_identity(), 1823 Some(generated.identity().public_key()) 1824 ); 1825 1826 let missing = 1827 PublicKey::from_hex("79be667ef9dcbbac55a06295ce870b07029bfcdb2dce28d959f2815b16f81798") 1828 .expect("public key"); 1829 let error = match actor.request_identity_removal(missing).await { 1830 Ok(_) => panic!("unknown identity removal must fail"), 1831 Err(error) => error, 1832 }; 1833 assert_eq!(error.code(), SafeErrorCode::IdentityNotFound); 1834 } 1835 1836 #[tokio::test(flavor = "multi_thread")] 1837 async fn staged_recovery_rejects_a_stale_expected_revision_before_commit() { 1838 let (actor, _) = actor().await; 1839 let handle = actor 1840 .begin_generated_key_stage() 1841 .await 1842 .expect("generated key stage"); 1843 let stale = SnapshotRevision::from_value(actor.snapshot().revision().value() + 1); 1844 let error = actor 1845 .acknowledge_generated_key_stage( 1846 handle.id(), 1847 DurableRequestId::parse("01890f3e-7b1c-7000-8000-000000005002").expect("request"), 1848 stale, 1849 DEFAULT_COMMAND_TIMEOUT, 1850 ) 1851 .await 1852 .expect_err("stale revision must conflict"); 1853 assert_eq!( 1854 error.message().as_str(), 1855 "The command conflicts with newer application state." 1856 ); 1857 assert!(actor.cancel_generated_key_stage().await.expect("cancel")); 1858 } 1859 1860 #[tokio::test(flavor = "multi_thread")] 1861 async fn fatal_lifecycle_rejects_commands_before_execution() { 1862 let (actor, _) = actor().await; 1863 actor 1864 .lifecycle 1865 .lock() 1866 .unwrap_or_else(std::sync::PoisonError::into_inner) 1867 .fail(command_unavailable()); 1868 1869 let error = actor 1870 .bootstrap() 1871 .await 1872 .expect_err("fatal lifecycle must reject command admission"); 1873 assert_eq!( 1874 error.message().as_str(), 1875 "The command is unavailable in the current runtime state." 1876 ); 1877 assert!(matches!(actor.lifecycle(), RuntimeLifecycle::Fatal(_))); 1878 } 1879 1880 #[tokio::test(flavor = "multi_thread")] 1881 async fn generated_key_stage_is_exclusive_cancelable_and_snapshot_free() { 1882 let (actor, secrets) = actor().await; 1883 let initial = actor.snapshot(); 1884 let stage = actor 1885 .begin_generated_key_stage() 1886 .await 1887 .expect("generated key stage"); 1888 1889 assert!(actor.begin_generated_key_stage().await.is_err()); 1890 assert_eq!(actor.snapshot(), initial); 1891 assert!( 1892 !secrets 1893 .contains(stage.view().identity().public_key()) 1894 .await 1895 .expect("keyring") 1896 ); 1897 assert!(actor.sign_out().await.is_err()); 1898 assert_eq!(actor.snapshot(), initial); 1899 assert!(actor.cancel_generated_key_stage().await.expect("cancel")); 1900 assert!( 1901 !actor 1902 .cancel_generated_key_stage() 1903 .await 1904 .expect("cancel empty") 1905 ); 1906 assert_eq!(actor.snapshot(), initial); 1907 1908 actor 1909 .begin_generated_key_stage() 1910 .await 1911 .expect("replacement stage"); 1912 actor.close().await.expect("close clears stage"); 1913 assert_eq!(actor.lifecycle(), RuntimeLifecycle::Closed); 1914 } 1915 1916 #[tokio::test(flavor = "multi_thread")] 1917 async fn recovery_handle_is_one_use_and_acknowledgement_commits_once() { 1918 let (actor, secrets) = actor().await; 1919 let initial = actor.snapshot(); 1920 let handle = actor 1921 .begin_generated_key_stage() 1922 .await 1923 .expect("generated key stage"); 1924 let public_key = handle.view().identity().public_key(); 1925 let recovery = handle.take_recovery_nsec().expect("recovery material"); 1926 assert_eq!(recovery.with_exposed_secret(str::len), 63); 1927 assert!(handle.take_recovery_nsec().is_err()); 1928 assert_eq!(actor.snapshot(), initial); 1929 assert!(!secrets.contains(public_key).await.expect("not committed")); 1930 1931 let committed = actor 1932 .acknowledge_generated_key_stage_test(handle.id()) 1933 .await 1934 .expect("acknowledge"); 1935 assert_eq!(committed.identities().len(), 1); 1936 assert_eq!(committed.selected_identity(), Some(public_key)); 1937 assert!( 1938 secrets 1939 .contains(public_key) 1940 .await 1941 .expect("credential committed") 1942 ); 1943 assert!( 1944 actor 1945 .acknowledge_generated_key_stage_test(handle.id()) 1946 .await 1947 .is_err() 1948 ); 1949 } 1950 1951 #[tokio::test(flavor = "multi_thread")] 1952 async fn failed_generated_commit_consumes_the_stage_without_poisoning_the_actor() { 1953 let secrets = Arc::new(FailureSecretStore::default()); 1954 secrets.fail_next(SecretStoreOperation::Put); 1955 let secret_port: Arc<dyn SecretStore> = secrets.clone(); 1956 let actor = RuntimeActorHandle::in_memory( 1957 RelayConfiguration::default(), 1958 dependencies(secret_port, Arc::new(OfflineNostr)), 1959 NonZeroUsize::new(8).expect("capacity"), 1960 &tokio::runtime::Handle::current(), 1961 ) 1962 .await 1963 .expect("actor"); 1964 let handle = actor 1965 .begin_generated_key_stage() 1966 .await 1967 .expect("generated key stage"); 1968 1969 let error = actor 1970 .acknowledge_generated_key_stage_test(handle.id()) 1971 .await 1972 .expect_err("injected keyring failure"); 1973 1974 assert_eq!(error.code(), SafeErrorCode::KeyringUnavailable); 1975 assert!(actor.snapshot().identities().is_empty()); 1976 actor 1977 .begin_generated_key_stage() 1978 .await 1979 .expect("fresh recovery after terminal failure"); 1980 assert!(actor.cancel_generated_key_stage().await.expect("cancel")); 1981 } 1982 1983 #[tokio::test(flavor = "multi_thread")] 1984 async fn session_generation_cancels_correlated_profile_work_on_sign_out() { 1985 let client = Arc::new(BlockingNostr::new()); 1986 let actor = RuntimeActorHandle::in_memory( 1987 RelayConfiguration::new(vec![ 1988 RelayEndpoint::new( 1989 "ws://localhost:8080", 1990 RelayUrlPolicy::Local, 1991 RelayAccess::ReadWrite, 1992 ) 1993 .expect("relay"), 1994 ]) 1995 .expect("relay configuration"), 1996 dependencies(Arc::new(InMemorySecretStore::default()), client.clone()), 1997 NonZeroUsize::new(8).expect("capacity"), 1998 &tokio::runtime::Handle::current(), 1999 ) 2000 .await 2001 .expect("actor"); 2002 let imported = actor 2003 .import_secret_key_test( 2004 SecretKeyInput::parse( 2005 "7e7e9c42a91bfef19fa7ea99d52d8afdb67d893a8fefba1f5cb9793f2107f6d7".to_owned(), 2006 ) 2007 .expect("input"), 2008 ) 2009 .await 2010 .expect("import"); 2011 actor 2012 .activate_identity(imported.identity().public_key()) 2013 .await 2014 .expect("activate"); 2015 assert_eq!(actor.session_generation().value(), 1); 2016 2017 let refresh_actor = actor.clone(); 2018 let refresh = tokio::spawn(async move { refresh_actor.refresh_active_profile().await }); 2019 let started = client.started.acquire().await.expect("refresh started"); 2020 started.forget(); 2021 let signed_out = actor.sign_out().await.expect("sign out"); 2022 let cancelled = refresh 2023 .await 2024 .expect("refresh task") 2025 .expect("safe cancellation"); 2026 2027 assert_eq!(actor.session_generation().value(), 2); 2028 assert_eq!(signed_out.session(), SessionState::SignedOut); 2029 assert_eq!(cancelled.session(), SessionState::SignedOut); 2030 assert!(cancelled.active_identity().is_none()); 2031 } 2032 2033 #[tokio::test(flavor = "multi_thread", worker_threads = 4)] 2034 async fn profile_refresh_rejects_stale_bindings_and_discards_stale_completions() { 2035 let client = Arc::new(BlockingNostr::new()); 2036 let actor = RuntimeActorHandle::in_memory( 2037 RelayConfiguration::new(vec![ 2038 RelayEndpoint::new( 2039 "ws://localhost:8080", 2040 RelayUrlPolicy::Local, 2041 RelayAccess::ReadWrite, 2042 ) 2043 .expect("relay"), 2044 ]) 2045 .expect("relay configuration"), 2046 dependencies(Arc::new(InMemorySecretStore::default()), client.clone()), 2047 NonZeroUsize::new(8).expect("capacity"), 2048 &tokio::runtime::Handle::current(), 2049 ) 2050 .await 2051 .expect("actor"); 2052 let imported = actor 2053 .import_secret_key_test(secret( 2054 "7e7e9c42a91bfef19fa7ea99d52d8afdb67d893a8fefba1f5cb9793f2107f6d7", 2055 )) 2056 .await 2057 .expect("import"); 2058 let public_key = imported.identity().public_key(); 2059 actor.activate_identity(public_key).await.expect("activate"); 2060 let binding = actor.foreground_session().expect("foreground binding"); 2061 let stale_binding = ActiveSessionBinding::new( 2062 NostrIdentityReference::derive(public_key).expect("identity"), 2063 LocalKeyringBinding::new(public_key, SignerAvailability::Available), 2064 SessionGeneration::from_value(binding.generation().value() + 1), 2065 ) 2066 .expect("stale binding fixture"); 2067 let other_public_key = 2068 PublicKey::from_hex("c6047f9441ed7d6d3045406e95c07cd85c778e4b8cef3ca7abac09b95c709ee5") 2069 .expect("other public key"); 2070 let other_binding = ActiveSessionBinding::new( 2071 NostrIdentityReference::derive(other_public_key).expect("other identity"), 2072 LocalKeyringBinding::new(other_public_key, SignerAvailability::Available), 2073 binding.generation(), 2074 ) 2075 .expect("other binding fixture"); 2076 2077 *actor 2078 .foreground_session 2079 .lock() 2080 .unwrap_or_else(std::sync::PoisonError::into_inner) = Some(stale_binding.clone()); 2081 let error = actor 2082 .refresh_active_profile() 2083 .await 2084 .expect_err("stale generation must reject before relay work"); 2085 assert_eq!( 2086 error.message().as_str(), 2087 "The active identity binding changed before profile refresh." 2088 ); 2089 2090 *actor 2091 .foreground_session 2092 .lock() 2093 .unwrap_or_else(std::sync::PoisonError::into_inner) = Some(other_binding.clone()); 2094 let error = actor 2095 .refresh_active_profile() 2096 .await 2097 .expect_err("different identity binding must reject before relay work"); 2098 assert_eq!( 2099 error.message().as_str(), 2100 "The active identity binding changed before profile refresh." 2101 ); 2102 2103 *actor 2104 .foreground_session 2105 .lock() 2106 .unwrap_or_else(std::sync::PoisonError::into_inner) = Some(binding.clone()); 2107 let refresh_actor = actor.clone(); 2108 let refresh = tokio::spawn(async move { refresh_actor.refresh_active_profile().await }); 2109 let started = client.started.acquire().await.expect("refresh started"); 2110 started.forget(); 2111 *actor 2112 .foreground_session 2113 .lock() 2114 .unwrap_or_else(std::sync::PoisonError::into_inner) = Some(stale_binding); 2115 client.release.add_permits(1); 2116 let unchanged = refresh 2117 .await 2118 .expect("refresh task") 2119 .expect("stale completion returns current snapshot"); 2120 assert_eq!(unchanged.revision(), actor.snapshot().revision()); 2121 2122 *actor 2123 .foreground_session 2124 .lock() 2125 .unwrap_or_else(std::sync::PoisonError::into_inner) = Some(binding.clone()); 2126 let refresh_actor = actor.clone(); 2127 let refresh = tokio::spawn(async move { refresh_actor.refresh_active_profile().await }); 2128 let started = client 2129 .started 2130 .acquire() 2131 .await 2132 .expect("second refresh started"); 2133 started.forget(); 2134 *actor 2135 .foreground_session 2136 .lock() 2137 .unwrap_or_else(std::sync::PoisonError::into_inner) = None; 2138 client.release.add_permits(1); 2139 let unchanged = refresh 2140 .await 2141 .expect("refresh task") 2142 .expect("missing binding returns current snapshot"); 2143 assert_eq!(unchanged.revision(), actor.snapshot().revision()); 2144 2145 *actor 2146 .foreground_session 2147 .lock() 2148 .unwrap_or_else(std::sync::PoisonError::into_inner) = Some(binding.clone()); 2149 let refresh_actor = actor.clone(); 2150 let refresh = tokio::spawn(async move { refresh_actor.refresh_active_profile().await }); 2151 let started = client 2152 .started 2153 .acquire() 2154 .await 2155 .expect("third refresh started"); 2156 started.forget(); 2157 *actor 2158 .foreground_session 2159 .lock() 2160 .unwrap_or_else(std::sync::PoisonError::into_inner) = Some(other_binding); 2161 client.release.add_permits(1); 2162 let unchanged = refresh 2163 .await 2164 .expect("refresh task") 2165 .expect("different identity binding returns current snapshot"); 2166 assert_eq!(unchanged.revision(), actor.snapshot().revision()); 2167 2168 *actor 2169 .foreground_session 2170 .lock() 2171 .unwrap_or_else(std::sync::PoisonError::into_inner) = Some(binding); 2172 } 2173 2174 #[tokio::test(flavor = "multi_thread", worker_threads = 4)] 2175 async fn bounded_runtime_rejects_saturation_without_dropping_accepted_commands() { 2176 let secrets = Arc::new(BlockingSecretStore::new()); 2177 let actor = RuntimeActorHandle::in_memory( 2178 RelayConfiguration::default(), 2179 dependencies(secrets.clone(), Arc::new(OfflineNostr)), 2180 NonZeroUsize::new(1).expect("capacity"), 2181 &tokio::runtime::Handle::current(), 2182 ) 2183 .await 2184 .expect("actor"); 2185 2186 let first_actor = actor.clone(); 2187 let first = tokio::spawn(async move { 2188 first_actor 2189 .import_secret_key_test(secret( 2190 "7e7e9c42a91bfef19fa7ea99d52d8afdb67d893a8fefba1f5cb9793f2107f6d7", 2191 )) 2192 .await 2193 }); 2194 secrets.wait_until_put_started().await; 2195 2196 let second_actor = actor.clone(); 2197 let second = tokio::spawn(async move { 2198 second_actor 2199 .import_secret_key_test(secret( 2200 "0000000000000000000000000000000000000000000000000000000000000001", 2201 )) 2202 .await 2203 }); 2204 while actor.mailbox.available_capacity() != 0 { 2205 assert!( 2206 !second.is_finished(), 2207 "second command must enter the mailbox" 2208 ); 2209 tokio::task::yield_now().await; 2210 } 2211 let rejected = actor 2212 .import_secret_key_test(secret( 2213 "0000000000000000000000000000000000000000000000000000000000000002", 2214 )) 2215 .await 2216 .expect_err("full mailbox must reject"); 2217 assert_eq!( 2218 rejected.message().as_str(), 2219 "The runtime is busy. Try again." 2220 ); 2221 2222 secrets.release(); 2223 first.await.expect("first task").expect("first command"); 2224 let second = second 2225 .await 2226 .expect("second task") 2227 .expect_err("accepted stale revision conflicts explicitly"); 2228 assert_eq!( 2229 second.message().as_str(), 2230 "The identity operation conflicts with the current application state." 2231 ); 2232 assert_eq!( 2233 actor 2234 .bootstrap() 2235 .await 2236 .expect("snapshot") 2237 .identities() 2238 .len(), 2239 1 2240 ); 2241 } 2242 2243 #[tokio::test(flavor = "multi_thread", worker_threads = 4)] 2244 async fn queued_command_expiry_returns_timeout_and_prevents_late_mutation() { 2245 let secrets = Arc::new(BlockingSecretStore::new()); 2246 let actor = RuntimeActorHandle::in_memory( 2247 RelayConfiguration::default(), 2248 dependencies(secrets.clone(), Arc::new(OfflineNostr)), 2249 NonZeroUsize::new(1).expect("capacity"), 2250 &tokio::runtime::Handle::current(), 2251 ) 2252 .await 2253 .expect("actor"); 2254 2255 let first_actor = actor.clone(); 2256 let first = tokio::spawn(async move { 2257 first_actor 2258 .import_secret_key_test(secret( 2259 "7e7e9c42a91bfef19fa7ea99d52d8afdb67d893a8fefba1f5cb9793f2107f6d7", 2260 )) 2261 .await 2262 }); 2263 secrets.wait_until_put_started().await; 2264 2265 let expired = actor 2266 .import_secret_key_with_timeout( 2267 secret("0000000000000000000000000000000000000000000000000000000000000001"), 2268 Duration::from_millis(10), 2269 ) 2270 .await 2271 .expect_err("queued command must time out"); 2272 assert_eq!(expired.message().as_str(), "The runtime command timed out."); 2273 2274 secrets.release(); 2275 first.await.expect("first task").expect("first command"); 2276 assert_eq!( 2277 actor 2278 .bootstrap() 2279 .await 2280 .expect("snapshot") 2281 .identities() 2282 .len(), 2283 1 2284 ); 2285 } 2286 2287 #[tokio::test(flavor = "multi_thread")] 2288 async fn close_is_terminal_and_every_later_command_is_rejected_as_closed() { 2289 let (actor, _) = actor().await; 2290 actor.close().await.expect("close"); 2291 assert_eq!(actor.lifecycle(), RuntimeLifecycle::Closed); 2292 2293 let error = actor.bootstrap().await.expect_err("bootstrap after close"); 2294 assert_eq!( 2295 error.message().as_str(), 2296 "The application runtime is closed." 2297 ); 2298 actor.close().await.expect("repeated close is idempotent"); 2299 } 2300 2301 #[tokio::test(flavor = "multi_thread")] 2302 async fn actor_subscription_atomically_delivers_initial_then_ordered_changes() { 2303 let (actor, _) = actor().await; 2304 let mut subscription = actor 2305 .subscribe_changes(NonZeroUsize::new(4).expect("capacity")) 2306 .await 2307 .expect("subscribe"); 2308 let initial = subscription.receive().await.expect("initial snapshot"); 2309 assert_eq!(initial.revision(), actor.snapshot().revision()); 2310 assert!(initial.previous_revision().is_none()); 2311 2312 actor 2313 .import_secret_key_test(secret( 2314 "7e7e9c42a91bfef19fa7ea99d52d8afdb67d893a8fefba1f5cb9793f2107f6d7", 2315 )) 2316 .await 2317 .expect("import"); 2318 let changed = subscription.receive().await.expect("change"); 2319 assert!(changed.revision() > initial.revision()); 2320 assert_eq!(changed.previous_revision(), Some(initial.revision())); 2321 assert!( 2322 actor 2323 .unsubscribe_changes(subscription.id()) 2324 .await 2325 .expect("unsubscribe") 2326 ); 2327 assert!( 2328 !actor 2329 .unsubscribe_changes(subscription.id()) 2330 .await 2331 .expect("second unsubscribe") 2332 ); 2333 } 2334 2335 #[tokio::test] 2336 async fn cancelled_registration_before_installation_does_not_retain_actor_admission() { 2337 let (handle, secrets) = actor().await; 2338 let snapshot = handle.snapshot(); 2339 // Obtain the first identifier from an independent probe, without relying on its value. 2340 let mut probe = OrderedSnapshotChanges::new(snapshot.clone()); 2341 let (probe_id, probe_receiver) = probe.subscribe(NonZeroUsize::MIN).expect("probe id"); 2342 drop(probe_receiver); 2343 let mut owner = RuntimeActor { 2344 adapter: Arc::clone(&handle.adapter), 2345 secrets, 2346 clock: Arc::new(FixedClock), 2347 nostr: Arc::new(OfflineNostr), 2348 lifecycle: Arc::clone(&handle.lifecycle), 2349 runtime: tokio::runtime::Handle::current(), 2350 session_generation: handle.session_generation(), 2351 published_session_generation: Arc::clone(&handle.session_generation), 2352 profile_tasks: std::collections::BTreeMap::new(), 2353 changes: OrderedSnapshotChanges::new(snapshot), 2354 published_foreground_session: Arc::clone(&handle.foreground_session), 2355 generated_key_stage: GeneratedKeyStage::default(), 2356 }; 2357 let (mailbox, mut receiver) = ActorMailbox::bounded(NonZeroUsize::MIN); 2358 let context = CommandContext::new( 2359 RequestId::new(1).expect("request id"), 2360 None, 2361 Instant::now() + DEFAULT_COMMAND_TIMEOUT, 2362 ); 2363 let ticket = match mailbox 2364 .submit(context, RuntimeCommand::SubscribeChanges(NonZeroUsize::MIN)) 2365 { 2366 CommandSubmission::Accepted(ticket) => ticket, 2367 CommandSubmission::Rejected(_) => panic!("empty test mailbox must admit registration"), 2368 }; 2369 drop(ticket); 2370 let envelope = receiver.recv().await.expect("queued registration"); 2371 let (completions, _completed) = tokio::sync::mpsc::channel(1); 2372 assert!(owner.handle_command(envelope, &completions).await); 2373 let retained_abandoned_registration = owner.changes.unsubscribe(probe_id); 2374 drop(owner); 2375 handle.close().await.expect("cleanup runtime"); 2376 2377 assert!( 2378 !retained_abandoned_registration, 2379 "cancelled queued registration must not retain admission" 2380 ); 2381 } 2382 2383 #[tokio::test] 2384 async fn abandoned_installed_subscription_releases_actor_admission_without_mutation() { 2385 let (actor, _) = actor().await; 2386 let subscription = actor 2387 .subscribe_changes(NonZeroUsize::MIN) 2388 .await 2389 .expect("installed actor subscription"); 2390 let id = subscription.id(); 2391 let unchanged_revision = actor.snapshot().revision(); 2392 drop(subscription); 2393 // Awaiting an Observe command orders any cleanup through the same actor mailbox. 2394 let snapshot = actor.bootstrap().await.expect("serialized cleanup barrier"); 2395 let retained_abandoned_registration = actor 2396 .unsubscribe_changes(id) 2397 .await 2398 .expect("registration probe"); 2399 actor.close().await.expect("cleanup runtime"); 2400 2401 assert_eq!(snapshot.revision(), unchanged_revision); 2402 assert!( 2403 !retained_abandoned_registration, 2404 "receiver abandonment must release admission without a later state mutation" 2405 ); 2406 } 2407 2408 #[tokio::test] 2409 async fn full_actor_queues_preserve_initial_and_final_tail_then_terminate_on_close() { 2410 let (actor, _) = actor().await; 2411 let mut public_keys = Vec::new(); 2412 for value in [ 2413 "7e7e9c42a91bfef19fa7ea99d52d8afdb67d893a8fefba1f5cb9793f2107f6d7", 2414 "6e7e9c42a91bfef19fa7ea99d52d8afdb67d893a8fefba1f5cb9793f2107f6d7", 2415 ] { 2416 public_keys.push( 2417 actor 2418 .import_secret_key_test(secret(value)) 2419 .await 2420 .expect("memory-store import") 2421 .identity() 2422 .public_key(), 2423 ); 2424 } 2425 actor 2426 .select_identity(public_keys[1]) 2427 .await 2428 .expect("initial selection"); 2429 let initial_revision = actor.snapshot().revision(); 2430 let mut first = actor 2431 .subscribe_changes(NonZeroUsize::MIN) 2432 .await 2433 .expect("first consumer"); 2434 let mut second = actor 2435 .subscribe_changes(NonZeroUsize::MIN) 2436 .await 2437 .expect("second consumer"); 2438 for offset in 0..4 { 2439 let selected = actor 2440 .select_identity(public_keys[offset % 2]) 2441 .await 2442 .expect("revision-changing selection"); 2443 assert_eq!( 2444 selected.revision().value(), 2445 initial_revision.value() + u64::try_from(offset).expect("offset") + 1 2446 ); 2447 } 2448 let final_snapshot = actor.snapshot(); 2449 tokio::time::timeout(Duration::from_secs(1), actor.close()) 2450 .await 2451 .expect("bounded full-queue close") 2452 .expect("close"); 2453 actor.close().await.expect("repeated close"); 2454 for subscription in [&mut first, &mut second] { 2455 let initial = tokio::time::timeout(Duration::from_secs(1), subscription.receive()) 2456 .await 2457 .expect("initial deadline") 2458 .expect("initial"); 2459 assert_eq!(initial.revision(), initial_revision); 2460 assert!(initial.previous_revision().is_none()); 2461 let tail = tokio::time::timeout(Duration::from_secs(1), subscription.receive()) 2462 .await 2463 .expect("tail deadline") 2464 .expect("final tail"); 2465 assert_eq!(tail.snapshot(), &final_snapshot); 2466 assert_eq!( 2467 tail.previous_revision() 2468 .expect("producer predecessor") 2469 .value(), 2470 final_snapshot.revision().value() - 1 2471 ); 2472 for _ in 0..2 { 2473 assert!( 2474 tokio::time::timeout(Duration::from_secs(1), subscription.receive()) 2475 .await 2476 .expect("terminal delivery deadline") 2477 .is_none() 2478 ); 2479 } 2480 } 2481 assert_eq!(actor.lifecycle(), RuntimeLifecycle::Closed); 2482 } 2483 2484 #[tokio::test(flavor = "multi_thread", worker_threads = 4)] 2485 async fn expired_queued_shutdown_does_not_close_runtime_later() { 2486 let secrets = Arc::new(BlockingSecretStore::new()); 2487 let actor = RuntimeActorHandle::in_memory( 2488 RelayConfiguration::default(), 2489 dependencies(secrets.clone(), Arc::new(OfflineNostr)), 2490 NonZeroUsize::new(1).expect("capacity"), 2491 &tokio::runtime::Handle::current(), 2492 ) 2493 .await 2494 .expect("actor"); 2495 let import_actor = actor.clone(); 2496 let import = tokio::spawn(async move { 2497 import_actor 2498 .import_secret_key_test(secret( 2499 "7e7e9c42a91bfef19fa7ea99d52d8afdb67d893a8fefba1f5cb9793f2107f6d7", 2500 )) 2501 .await 2502 }); 2503 secrets.wait_until_put_started().await; 2504 2505 let timeout = actor 2506 .close_with_timeout(Duration::from_millis(10)) 2507 .await 2508 .expect_err("queued shutdown must expire"); 2509 assert_eq!(timeout.message().as_str(), "The runtime command timed out."); 2510 secrets.release(); 2511 import.await.expect("import task").expect("import"); 2512 assert_eq!(actor.lifecycle(), RuntimeLifecycle::Ready); 2513 assert_eq!( 2514 actor 2515 .bootstrap() 2516 .await 2517 .expect("still open") 2518 .identities() 2519 .len(), 2520 1 2521 ); 2522 actor.close().await.expect("later close"); 2523 } 2524 2525 #[tokio::test(flavor = "multi_thread", worker_threads = 4)] 2526 async fn shutdown_cancels_in_flight_work_and_terminates_publication() { 2527 let client = Arc::new(BlockingNostr::new()); 2528 let actor = RuntimeActorHandle::in_memory( 2529 RelayConfiguration::new(vec![ 2530 RelayEndpoint::new( 2531 "ws://localhost:8080", 2532 RelayUrlPolicy::Local, 2533 RelayAccess::ReadWrite, 2534 ) 2535 .expect("relay"), 2536 ]) 2537 .expect("relay configuration"), 2538 dependencies(Arc::new(InMemorySecretStore::default()), client.clone()), 2539 NonZeroUsize::new(8).expect("capacity"), 2540 &tokio::runtime::Handle::current(), 2541 ) 2542 .await 2543 .expect("actor"); 2544 let imported = actor 2545 .import_secret_key_test(secret( 2546 "7e7e9c42a91bfef19fa7ea99d52d8afdb67d893a8fefba1f5cb9793f2107f6d7", 2547 )) 2548 .await 2549 .expect("import"); 2550 actor 2551 .activate_identity(imported.identity().public_key()) 2552 .await 2553 .expect("activate"); 2554 let mut changes = actor 2555 .subscribe_changes(NonZeroUsize::new(4).expect("capacity")) 2556 .await 2557 .expect("subscribe"); 2558 changes.receive().await.expect("initial"); 2559 2560 let refresh_actor = actor.clone(); 2561 let refresh = tokio::spawn(async move { refresh_actor.refresh_active_profile().await }); 2562 let started = client.started.acquire().await.expect("refresh started"); 2563 started.forget(); 2564 actor.close().await.expect("close"); 2565 2566 let cancelled = refresh 2567 .await 2568 .expect("refresh task") 2569 .expect_err("refresh closes"); 2570 assert_eq!( 2571 cancelled.message().as_str(), 2572 "The application runtime is closed." 2573 ); 2574 assert!(changes.receive().await.is_none()); 2575 assert_eq!(actor.lifecycle(), RuntimeLifecycle::Closed); 2576 } 2577 2578 fn secret(value: &str) -> SecretKeyInput { 2579 SecretKeyInput::parse(value.to_owned()).expect("valid test secret") 2580 } 2581 }