recovery.rs (33092B)
1 #[cfg(test)] 2 use harvestcircle_domain::PublicKey; 3 use harvestcircle_domain::SafeError; 4 5 use crate::{ 6 AppCore, AppStateRepository, Clock, DurableIdentityOperation, DurableOperationKind, 7 DurableOperationPhase, DurableOperationRepository, DurableTerminalOutcome, IdentityRepository, 8 SecretStore, 9 }; 10 #[cfg(test)] 11 use crate::{DurableRequestId, IdentityOperationKind, IdentityOperationPhase, OperationJournal}; 12 13 impl AppCore { 14 /// Reconciles durable request operations before public state is restored. 15 /// 16 /// # Errors 17 /// 18 /// Returns a safe credential, persistence, or recovery error while retaining the operation 19 /// at its last durable phase for a later retry. 20 pub async fn recover_durable_operations( 21 &self, 22 identities: &(impl IdentityRepository + ?Sized), 23 app_state: &(impl AppStateRepository + ?Sized), 24 secrets: &(impl SecretStore + ?Sized), 25 operations: &(impl DurableOperationRepository + ?Sized), 26 clock: &(impl Clock + ?Sized), 27 ) -> Result<(), SafeError> { 28 for operation in operations.list_unfinished_durable_operations().await? { 29 match operation.kind() { 30 DurableOperationKind::Create 31 | DurableOperationKind::Import 32 | DurableOperationKind::Repair => { 33 recover_durable_addition( 34 &operation, identities, app_state, secrets, operations, clock, 35 ) 36 .await? 37 } 38 DurableOperationKind::Remove => { 39 recover_durable_removal( 40 &operation, identities, app_state, secrets, operations, clock, 41 ) 42 .await? 43 } 44 } 45 } 46 Ok(()) 47 } 48 49 /// Reconciles non-secret cross-resource journal entries before bootstrap. 50 /// 51 /// An empty journal does not access the credential store. 52 /// 53 /// # Errors 54 /// 55 /// Returns a safe credential, persistence, or recovery error while retaining 56 /// the unfinished journal entry for a later retry. 57 #[cfg(test)] 58 pub async fn recover_pending_operations( 59 &self, 60 identities: &(impl IdentityRepository + ?Sized), 61 app_state: &(impl AppStateRepository + ?Sized), 62 secrets: &(impl SecretStore + ?Sized), 63 journal: &(impl OperationJournal + ?Sized), 64 clock: &(impl Clock + ?Sized), 65 ) -> Result<(), SafeError> { 66 for operation in journal.list_pending_operations().await? { 67 match operation.kind() { 68 IdentityOperationKind::Remove => { 69 recover_removal(&operation, identities, app_state, secrets, journal, clock) 70 .await?; 71 } 72 IdentityOperationKind::Add | IdentityOperationKind::Import => { 73 recover_addition(&operation, identities, secrets, journal, clock).await?; 74 } 75 } 76 } 77 Ok(()) 78 } 79 } 80 81 async fn recover_durable_removal( 82 operation: &DurableIdentityOperation, 83 identities: &(impl IdentityRepository + ?Sized), 84 app_state: &(impl AppStateRepository + ?Sized), 85 secrets: &(impl SecretStore + ?Sized), 86 operations: &(impl DurableOperationRepository + ?Sized), 87 clock: &(impl Clock + ?Sized), 88 ) -> Result<(), SafeError> { 89 let request = operation.request_id(); 90 let identity = operation.identity(); 91 let mut phase = operation.phase(); 92 if phase == DurableOperationPhase::IntentRecorded { 93 if secrets.contains(identity).await? { 94 secrets.delete(operation.request_id(), identity).await?; 95 } 96 operations 97 .advance_durable_operation( 98 request, 99 phase, 100 DurableOperationPhase::CredentialDeleted, 101 clock.now(), 102 None, 103 ) 104 .await?; 105 phase = DurableOperationPhase::CredentialDeleted; 106 } 107 if phase == DurableOperationPhase::CredentialDeleted { 108 if identities.find_identity(identity).await?.is_some() { 109 identities.remove_identity(identity).await?; 110 } 111 operations 112 .advance_durable_operation( 113 request, 114 phase, 115 DurableOperationPhase::MetadataDeleted, 116 clock.now(), 117 None, 118 ) 119 .await?; 120 phase = DurableOperationPhase::MetadataDeleted; 121 } 122 if phase == DurableOperationPhase::MetadataDeleted { 123 app_state 124 .save_selected_identity(operation.prior().selected_identity()) 125 .await?; 126 operations 127 .advance_durable_operation( 128 request, 129 phase, 130 DurableOperationPhase::SelectionCommitted, 131 clock.now(), 132 None, 133 ) 134 .await?; 135 phase = DurableOperationPhase::SelectionCommitted; 136 } 137 if phase == DurableOperationPhase::SelectionCommitted { 138 operations 139 .finalize_durable_operation( 140 request, 141 phase, 142 DurableTerminalOutcome::Completed, 143 None, 144 clock.now(), 145 ) 146 .await?; 147 } 148 Ok(()) 149 } 150 151 async fn recover_durable_addition( 152 operation: &DurableIdentityOperation, 153 identities: &(impl IdentityRepository + ?Sized), 154 app_state: &(impl AppStateRepository + ?Sized), 155 secrets: &(impl SecretStore + ?Sized), 156 operations: &(impl DurableOperationRepository + ?Sized), 157 clock: &(impl Clock + ?Sized), 158 ) -> Result<(), SafeError> { 159 let request = operation.request_id(); 160 let identity = operation.identity(); 161 match operation.phase() { 162 DurableOperationPhase::IntentRecorded => { 163 if secrets.contains(identity).await? { 164 secrets.delete(operation.request_id(), identity).await?; 165 } 166 operations 167 .finalize_durable_operation( 168 request, 169 DurableOperationPhase::IntentRecorded, 170 DurableTerminalOutcome::Failed, 171 None, 172 clock.now(), 173 ) 174 .await?; 175 } 176 DurableOperationPhase::CredentialWritten => { 177 let metadata = identities.find_identity(identity).await?; 178 let committed = metadata.as_ref().is_some_and(|saved| { 179 saved 180 .signer_binding() 181 .as_local_keyring() 182 .is_some_and(|binding| { 183 binding.availability() 184 == harvestcircle_domain::SignerAvailability::Available 185 }) 186 }); 187 if committed { 188 operations 189 .advance_durable_operation( 190 request, 191 DurableOperationPhase::CredentialWritten, 192 DurableOperationPhase::MetadataCommitted, 193 clock.now(), 194 None, 195 ) 196 .await?; 197 finish_durable_selection(operation, app_state, operations, clock).await?; 198 } else { 199 operations 200 .advance_durable_operation( 201 request, 202 DurableOperationPhase::CredentialWritten, 203 DurableOperationPhase::CompensationPending, 204 clock.now(), 205 None, 206 ) 207 .await?; 208 compensate_durable_addition( 209 operation, identities, app_state, secrets, operations, clock, 210 ) 211 .await?; 212 } 213 } 214 DurableOperationPhase::MetadataCommitted => { 215 finish_durable_selection(operation, app_state, operations, clock).await?; 216 } 217 DurableOperationPhase::SelectionCommitted => { 218 operations 219 .finalize_durable_operation( 220 request, 221 DurableOperationPhase::SelectionCommitted, 222 DurableTerminalOutcome::Completed, 223 None, 224 clock.now(), 225 ) 226 .await?; 227 } 228 DurableOperationPhase::CompensationPending => { 229 compensate_durable_addition( 230 operation, identities, app_state, secrets, operations, clock, 231 ) 232 .await?; 233 } 234 DurableOperationPhase::CredentialDeleted | DurableOperationPhase::MetadataDeleted => { 235 operations 236 .finalize_durable_operation( 237 request, 238 operation.phase(), 239 DurableTerminalOutcome::Failed, 240 None, 241 clock.now(), 242 ) 243 .await?; 244 } 245 DurableOperationPhase::Finalized => {} 246 } 247 Ok(()) 248 } 249 250 async fn finish_durable_selection( 251 operation: &DurableIdentityOperation, 252 app_state: &(impl AppStateRepository + ?Sized), 253 operations: &(impl DurableOperationRepository + ?Sized), 254 clock: &(impl Clock + ?Sized), 255 ) -> Result<(), SafeError> { 256 app_state 257 .save_selected_identity(Some(operation.identity())) 258 .await?; 259 operations 260 .advance_durable_operation( 261 operation.request_id(), 262 DurableOperationPhase::MetadataCommitted, 263 DurableOperationPhase::SelectionCommitted, 264 clock.now(), 265 None, 266 ) 267 .await?; 268 operations 269 .finalize_durable_operation( 270 operation.request_id(), 271 DurableOperationPhase::SelectionCommitted, 272 DurableTerminalOutcome::Completed, 273 None, 274 clock.now(), 275 ) 276 .await?; 277 Ok(()) 278 } 279 280 async fn compensate_durable_addition( 281 operation: &DurableIdentityOperation, 282 identities: &(impl IdentityRepository + ?Sized), 283 app_state: &(impl AppStateRepository + ?Sized), 284 secrets: &(impl SecretStore + ?Sized), 285 operations: &(impl DurableOperationRepository + ?Sized), 286 clock: &(impl Clock + ?Sized), 287 ) -> Result<(), SafeError> { 288 if secrets.contains(operation.identity()).await? { 289 secrets 290 .delete(operation.request_id(), operation.identity()) 291 .await?; 292 } 293 if let Some(availability) = operation.prior().binding_availability() { 294 if let Some(previous) = identities.find_identity(operation.identity()).await? { 295 let restored = previous 296 .with_local_keyring_availability(availability) 297 .ok_or_else(recovery_required)?; 298 identities.update_identity(&restored).await?; 299 } 300 } else if identities 301 .find_identity(operation.identity()) 302 .await? 303 .is_some() 304 { 305 identities.remove_identity(operation.identity()).await?; 306 } 307 app_state 308 .save_selected_identity(operation.prior().selected_identity()) 309 .await?; 310 operations 311 .advance_durable_operation( 312 operation.request_id(), 313 DurableOperationPhase::CompensationPending, 314 DurableOperationPhase::CredentialDeleted, 315 clock.now(), 316 None, 317 ) 318 .await?; 319 operations 320 .finalize_durable_operation( 321 operation.request_id(), 322 DurableOperationPhase::CredentialDeleted, 323 DurableTerminalOutcome::Failed, 324 None, 325 clock.now(), 326 ) 327 .await?; 328 Ok(()) 329 } 330 331 const fn recovery_required() -> SafeError { 332 SafeError::new( 333 harvestcircle_domain::SafeErrorCode::PendingOperationRecoveryRequired, 334 harvestcircle_domain::SafeMessage::new( 335 "Identity recovery is required before this operation can continue.", 336 ), 337 ) 338 } 339 340 #[cfg(test)] 341 async fn recover_removal( 342 operation: &crate::PendingIdentityOperation, 343 identities: &(impl IdentityRepository + ?Sized), 344 app_state: &(impl AppStateRepository + ?Sized), 345 secrets: &(impl SecretStore + ?Sized), 346 journal: &(impl OperationJournal + ?Sized), 347 clock: &(impl Clock + ?Sized), 348 ) -> Result<(), SafeError> { 349 let public_key = operation.subject(); 350 if operation.phase() == IdentityOperationPhase::IntentRecorded { 351 let request_id = DurableRequestId::new_v7(); 352 match secrets.delete(&request_id, public_key).await { 353 Ok(()) => {} 354 Err(error) 355 if error.code() == harvestcircle_domain::SafeErrorCode::CredentialMissing => {} 356 Err(error) => return Err(error), 357 } 358 journal 359 .update_operation( 360 operation.id(), 361 IdentityOperationPhase::CredentialDeleted, 362 clock.now(), 363 None, 364 ) 365 .await?; 366 } 367 if matches!( 368 operation.phase(), 369 IdentityOperationPhase::IntentRecorded | IdentityOperationPhase::CredentialDeleted 370 ) { 371 let registry = identities.list_identities().await?; 372 let selected = removal_fallback( 373 ®istry, 374 app_state.load_selected_identity().await?, 375 public_key, 376 ); 377 identities.remove_identity(public_key).await?; 378 app_state.save_selected_identity(selected).await?; 379 journal 380 .update_operation( 381 operation.id(), 382 IdentityOperationPhase::MetadataDeleted, 383 clock.now(), 384 None, 385 ) 386 .await?; 387 } 388 journal.finalize_operation(operation.id()).await 389 } 390 391 #[cfg(test)] 392 async fn recover_addition( 393 operation: &crate::PendingIdentityOperation, 394 identities: &(impl IdentityRepository + ?Sized), 395 secrets: &(impl SecretStore + ?Sized), 396 journal: &(impl OperationJournal + ?Sized), 397 clock: &(impl Clock + ?Sized), 398 ) -> Result<(), SafeError> { 399 let has_metadata = identities 400 .find_identity(operation.subject()) 401 .await? 402 .is_some(); 403 match operation.phase() { 404 IdentityOperationPhase::CredentialWritten | IdentityOperationPhase::CompensationPending 405 if !has_metadata => 406 { 407 let request_id = DurableRequestId::new_v7(); 408 match secrets.delete(&request_id, operation.subject()).await { 409 Ok(()) => {} 410 Err(error) 411 if error.code() == harvestcircle_domain::SafeErrorCode::CredentialMissing => {} 412 Err(error) => return Err(error), 413 } 414 journal 415 .update_operation( 416 operation.id(), 417 IdentityOperationPhase::MetadataDeleted, 418 clock.now(), 419 None, 420 ) 421 .await?; 422 } 423 _ => {} 424 } 425 journal.finalize_operation(operation.id()).await 426 } 427 428 #[cfg(test)] 429 fn removal_fallback( 430 registry: &[harvestcircle_domain::NostrIdentity], 431 selected: Option<PublicKey>, 432 removed: PublicKey, 433 ) -> Option<PublicKey> { 434 if selected != Some(removed) { 435 return selected; 436 } 437 let index = registry 438 .iter() 439 .position(|identity| identity.public_key() == removed)?; 440 registry 441 .get(index + 1) 442 .or_else(|| index.checked_sub(1).and_then(|before| registry.get(before))) 443 .map(harvestcircle_domain::NostrIdentity::public_key) 444 } 445 446 #[cfg(test)] 447 pub(crate) mod tests { 448 use std::sync::{Mutex, MutexGuard}; 449 450 use harvestcircle_domain::{ 451 PublicKey, SafeError, SafeErrorCode, SafeMessage, SecretKeyInput, SignerAvailability, 452 UnixTimestamp, 453 }; 454 455 use super::*; 456 use crate::{ 457 BoxFuture, DurableOperationReceipt, DurableOperationStart, DurableRequestId, 458 FailureSecretStore, InMemoryIdentityRepository, InMemoryOperationJournal, 459 InMemorySecretStore, RelayConfiguration, SecretStore, SecretStoreOperation, 460 }; 461 462 struct FixedClock; 463 464 impl Clock for FixedClock { 465 fn now(&self) -> UnixTimestamp { 466 UnixTimestamp::from_seconds(10).expect("time") 467 } 468 } 469 470 pub(crate) struct TestDurableRepository { 471 operation: Mutex<DurableIdentityOperation>, 472 return_existing: bool, 473 } 474 475 impl TestDurableRepository { 476 pub(crate) fn new(operation: DurableIdentityOperation) -> Self { 477 Self { 478 operation: Mutex::new(operation), 479 return_existing: true, 480 } 481 } 482 483 pub(crate) fn fresh(operation: DurableIdentityOperation) -> Self { 484 Self { 485 operation: Mutex::new(operation), 486 return_existing: false, 487 } 488 } 489 490 pub(crate) fn operation(&self) -> MutexGuard<'_, DurableIdentityOperation> { 491 self.operation 492 .lock() 493 .unwrap_or_else(std::sync::PoisonError::into_inner) 494 } 495 496 fn replace( 497 current: &DurableIdentityOperation, 498 phase: DurableOperationPhase, 499 diagnostic: Option<crate::OperationDiagnostic>, 500 terminal: Option<DurableOperationReceipt>, 501 ) -> DurableIdentityOperation { 502 DurableIdentityOperation::new( 503 current.request_id().clone(), 504 current.kind(), 505 current.identity(), 506 current.expected_revision(), 507 phase, 508 current.prior(), 509 current.updated_at(), 510 diagnostic, 511 terminal, 512 ) 513 } 514 } 515 516 impl DurableOperationRepository for TestDurableRepository { 517 fn begin_durable_operation<'a>( 518 &'a self, 519 _request_id: &'a DurableRequestId, 520 _kind: DurableOperationKind, 521 _identity: PublicKey, 522 _expected_revision: Option<u64>, 523 _prior: crate::OperationPriorState, 524 _updated_at: UnixTimestamp, 525 ) -> BoxFuture<'a, Result<DurableOperationStart, SafeError>> { 526 Box::pin(async move { 527 let operation = self.operation().clone(); 528 Ok(if self.return_existing { 529 DurableOperationStart::Existing(operation) 530 } else { 531 DurableOperationStart::Started(operation) 532 }) 533 }) 534 } 535 536 fn load_durable_operation<'a>( 537 &'a self, 538 request_id: &'a DurableRequestId, 539 ) -> BoxFuture<'a, Result<Option<DurableIdentityOperation>, SafeError>> { 540 Box::pin(async move { 541 if !self.return_existing { 542 return Ok(None); 543 } 544 let operation = self.operation(); 545 Ok((operation.request_id() == request_id).then(|| operation.clone())) 546 }) 547 } 548 549 fn advance_durable_operation<'a>( 550 &'a self, 551 request_id: &'a DurableRequestId, 552 expected_phase: DurableOperationPhase, 553 next_phase: DurableOperationPhase, 554 _updated_at: UnixTimestamp, 555 diagnostic: Option<crate::OperationDiagnostic>, 556 ) -> BoxFuture<'a, Result<DurableIdentityOperation, SafeError>> { 557 Box::pin(async move { 558 let mut operation = self.operation(); 559 if operation.request_id() != request_id || operation.phase() != expected_phase { 560 return Err(conflict()); 561 } 562 *operation = Self::replace(&operation, next_phase, diagnostic, None); 563 Ok(operation.clone()) 564 }) 565 } 566 567 fn finalize_durable_operation<'a>( 568 &'a self, 569 request_id: &'a DurableRequestId, 570 expected_phase: DurableOperationPhase, 571 outcome: DurableTerminalOutcome, 572 resulting_revision: Option<u64>, 573 updated_at: UnixTimestamp, 574 ) -> BoxFuture<'a, Result<DurableOperationReceipt, SafeError>> { 575 Box::pin(async move { 576 let mut operation = self.operation(); 577 if operation.request_id() != request_id || operation.phase() != expected_phase { 578 return Err(conflict()); 579 } 580 let receipt = DurableOperationReceipt::new( 581 request_id.clone(), 582 operation.identity(), 583 outcome, 584 resulting_revision, 585 updated_at, 586 ); 587 *operation = Self::replace( 588 &operation, 589 DurableOperationPhase::Finalized, 590 operation.diagnostic(), 591 Some(receipt.clone()), 592 ); 593 Ok(receipt) 594 }) 595 } 596 597 fn list_unfinished_durable_operations( 598 &self, 599 ) -> BoxFuture<'_, Result<Vec<DurableIdentityOperation>, SafeError>> { 600 Box::pin(async move { Ok(vec![self.operation().clone()]) }) 601 } 602 } 603 604 fn conflict() -> SafeError { 605 SafeError::new( 606 SafeErrorCode::InvalidApplicationState, 607 SafeMessage::new("The test durable operation conflicted."), 608 ) 609 } 610 611 async fn seeded() -> ( 612 AppCore, 613 InMemoryIdentityRepository, 614 InMemorySecretStore, 615 InMemoryOperationJournal, 616 PublicKey, 617 ) { 618 let core = AppCore::in_memory(RelayConfiguration::default()); 619 let identities = InMemoryIdentityRepository::default(); 620 let secrets = InMemorySecretStore::default(); 621 let journal = InMemoryOperationJournal::default(); 622 core.bootstrap().expect("bootstrap"); 623 let receipt = core 624 .generate_identity(&identities, &identities, &secrets, &journal, &FixedClock) 625 .await 626 .expect("seed identity"); 627 ( 628 core, 629 identities, 630 secrets, 631 journal, 632 receipt.identity().public_key(), 633 ) 634 } 635 636 pub(crate) fn operation( 637 kind: DurableOperationKind, 638 phase: DurableOperationPhase, 639 identity: PublicKey, 640 prior_availability: Option<SignerAvailability>, 641 ) -> DurableIdentityOperation { 642 DurableIdentityOperation::new( 643 DurableRequestId::new_v7(), 644 kind, 645 identity, 646 Some(1), 647 phase, 648 crate::OperationPriorState::new(None, prior_availability), 649 FixedClock.now(), 650 None, 651 None, 652 ) 653 } 654 655 async fn run_durable( 656 core: &AppCore, 657 identities: &InMemoryIdentityRepository, 658 secrets: &InMemorySecretStore, 659 operation: DurableIdentityOperation, 660 ) -> DurableIdentityOperation { 661 let repository = TestDurableRepository::new(operation); 662 core.recover_durable_operations(identities, identities, secrets, &repository, &FixedClock) 663 .await 664 .expect("durable recovery"); 665 repository.operation().clone() 666 } 667 668 #[tokio::test] 669 async fn durable_recovery_exercises_every_removal_phase_and_presence_branch() { 670 for phase in [ 671 DurableOperationPhase::IntentRecorded, 672 DurableOperationPhase::CredentialDeleted, 673 DurableOperationPhase::MetadataDeleted, 674 DurableOperationPhase::SelectionCommitted, 675 DurableOperationPhase::Finalized, 676 ] { 677 let (core, identities, secrets, _journal, public_key) = seeded().await; 678 if phase != DurableOperationPhase::IntentRecorded { 679 secrets 680 .delete(&DurableRequestId::new_v7(), public_key) 681 .await 682 .expect("delete credential"); 683 } 684 if matches!( 685 phase, 686 DurableOperationPhase::MetadataDeleted 687 | DurableOperationPhase::SelectionCommitted 688 | DurableOperationPhase::Finalized 689 ) { 690 identities 691 .remove_identity(public_key) 692 .await 693 .expect("remove identity"); 694 } 695 let recovered = run_durable( 696 &core, 697 &identities, 698 &secrets, 699 operation(DurableOperationKind::Remove, phase, public_key, None), 700 ) 701 .await; 702 assert_eq!(recovered.phase(), DurableOperationPhase::Finalized); 703 } 704 705 let (core, identities, secrets, _journal, public_key) = seeded().await; 706 secrets 707 .delete(&DurableRequestId::new_v7(), public_key) 708 .await 709 .expect("delete credential"); 710 identities 711 .remove_identity(public_key) 712 .await 713 .expect("remove identity"); 714 let recovered = run_durable( 715 &core, 716 &identities, 717 &secrets, 718 operation( 719 DurableOperationKind::Remove, 720 DurableOperationPhase::IntentRecorded, 721 public_key, 722 None, 723 ), 724 ) 725 .await; 726 assert_eq!(recovered.phase(), DurableOperationPhase::Finalized); 727 } 728 729 #[tokio::test] 730 async fn durable_recovery_exercises_every_addition_phase_and_compensation_shape() { 731 for phase in [ 732 DurableOperationPhase::IntentRecorded, 733 DurableOperationPhase::CredentialWritten, 734 DurableOperationPhase::MetadataCommitted, 735 DurableOperationPhase::SelectionCommitted, 736 DurableOperationPhase::CredentialDeleted, 737 DurableOperationPhase::MetadataDeleted, 738 DurableOperationPhase::Finalized, 739 ] { 740 let (core, identities, secrets, _journal, public_key) = seeded().await; 741 if phase == DurableOperationPhase::IntentRecorded { 742 identities 743 .remove_identity(public_key) 744 .await 745 .expect("remove metadata"); 746 } 747 let recovered = run_durable( 748 &core, 749 &identities, 750 &secrets, 751 operation(DurableOperationKind::Create, phase, public_key, None), 752 ) 753 .await; 754 assert_eq!(recovered.phase(), DurableOperationPhase::Finalized); 755 } 756 757 for (prior, retain_metadata, retain_secret) in [ 758 (Some(SignerAvailability::CredentialMissing), true, true), 759 (Some(SignerAvailability::CredentialMissing), false, true), 760 (None, true, true), 761 (None, false, false), 762 ] { 763 let (core, identities, secrets, _journal, public_key) = seeded().await; 764 if !retain_metadata { 765 identities 766 .remove_identity(public_key) 767 .await 768 .expect("remove metadata"); 769 } 770 if !retain_secret { 771 secrets 772 .delete(&DurableRequestId::new_v7(), public_key) 773 .await 774 .expect("delete credential"); 775 } 776 let recovered = run_durable( 777 &core, 778 &identities, 779 &secrets, 780 operation( 781 DurableOperationKind::Repair, 782 DurableOperationPhase::CompensationPending, 783 public_key, 784 prior, 785 ), 786 ) 787 .await; 788 assert_eq!(recovered.phase(), DurableOperationPhase::Finalized); 789 } 790 791 let (core, identities, secrets, _journal, public_key) = seeded().await; 792 identities 793 .remove_identity(public_key) 794 .await 795 .expect("remove metadata"); 796 let recovered = run_durable( 797 &core, 798 &identities, 799 &secrets, 800 operation( 801 DurableOperationKind::Import, 802 DurableOperationPhase::CredentialWritten, 803 public_key, 804 None, 805 ), 806 ) 807 .await; 808 assert_eq!(recovered.phase(), DurableOperationPhase::Finalized); 809 810 let (core, identities, secrets, _journal, public_key) = seeded().await; 811 identities 812 .remove_identity(public_key) 813 .await 814 .expect("remove metadata"); 815 secrets 816 .delete(&DurableRequestId::new_v7(), public_key) 817 .await 818 .expect("delete credential"); 819 let recovered = run_durable( 820 &core, 821 &identities, 822 &secrets, 823 operation( 824 DurableOperationKind::Create, 825 DurableOperationPhase::IntentRecorded, 826 public_key, 827 None, 828 ), 829 ) 830 .await; 831 assert_eq!(recovered.phase(), DurableOperationPhase::Finalized); 832 } 833 834 #[tokio::test] 835 async fn pending_recovery_exercises_removal_and_addition_presence_branches() { 836 for credential_present in [true, false] { 837 let (core, identities, secrets, journal, public_key) = seeded().await; 838 if !credential_present { 839 secrets 840 .delete(&DurableRequestId::new_v7(), public_key) 841 .await 842 .expect("delete credential"); 843 identities 844 .save_selected_identity(None) 845 .await 846 .expect("clear selection"); 847 } 848 journal 849 .begin_operation(IdentityOperationKind::Remove, public_key, FixedClock.now()) 850 .await 851 .expect("removal intent"); 852 core.recover_pending_operations( 853 &identities, 854 &identities, 855 &secrets, 856 &journal, 857 &FixedClock, 858 ) 859 .await 860 .expect("removal recovery"); 861 assert!(journal.list_pending_operations().await.unwrap().is_empty()); 862 } 863 864 for (kind, metadata_present, credential_present) in [ 865 (IdentityOperationKind::Add, false, true), 866 (IdentityOperationKind::Import, false, false), 867 (IdentityOperationKind::Add, true, true), 868 ] { 869 let (core, identities, secrets, journal, public_key) = seeded().await; 870 if !metadata_present { 871 identities 872 .remove_identity(public_key) 873 .await 874 .expect("remove metadata"); 875 } 876 if !credential_present { 877 secrets 878 .delete(&DurableRequestId::new_v7(), public_key) 879 .await 880 .expect("delete credential"); 881 } 882 let id = journal 883 .begin_operation(kind, public_key, FixedClock.now()) 884 .await 885 .expect("addition intent"); 886 journal 887 .update_operation( 888 id, 889 IdentityOperationPhase::CredentialWritten, 890 FixedClock.now(), 891 None, 892 ) 893 .await 894 .expect("credential phase"); 895 core.recover_pending_operations( 896 &identities, 897 &identities, 898 &secrets, 899 &journal, 900 &FixedClock, 901 ) 902 .await 903 .expect("addition recovery"); 904 assert!(journal.list_pending_operations().await.unwrap().is_empty()); 905 } 906 907 let (core, identities, _secrets, journal, public_key) = seeded().await; 908 identities 909 .remove_identity(public_key) 910 .await 911 .expect("remove metadata"); 912 let secrets = FailureSecretStore::default(); 913 secrets 914 .put( 915 &DurableRequestId::new_v7(), 916 public_key, 917 SecretKeyInput::parse( 918 "7e7e9c42a91bfef19fa7ea99d52d8afdb67d893a8fefba1f5cb9793f2107f6d7".to_owned(), 919 ) 920 .expect("secret"), 921 ) 922 .await 923 .expect("store credential"); 924 secrets.fail_next(SecretStoreOperation::Delete); 925 let id = journal 926 .begin_operation(IdentityOperationKind::Add, public_key, FixedClock.now()) 927 .await 928 .expect("addition intent"); 929 journal 930 .update_operation( 931 id, 932 IdentityOperationPhase::CredentialWritten, 933 FixedClock.now(), 934 None, 935 ) 936 .await 937 .expect("credential phase"); 938 assert!( 939 core.recover_pending_operations( 940 &identities, 941 &identities, 942 &secrets, 943 &journal, 944 &FixedClock, 945 ) 946 .await 947 .is_err() 948 ); 949 } 950 }