actor.rs (24343B)
1 use std::num::{NonZeroU64, NonZeroUsize}; 2 use std::time::Instant; 3 4 use harvestcircle_domain::{ 5 NostrIdentityReference, PublicKey, SafeError, SafeErrorCode, SafeMessage, SignerAvailability, 6 SignerBinding, 7 }; 8 use tokio::sync::{mpsc, oneshot}; 9 10 use crate::SnapshotRevision; 11 12 #[derive(Clone, Copy, Debug, Default, Eq, Hash, Ord, PartialEq, PartialOrd)] 13 pub struct SessionGeneration(u64); 14 15 impl SessionGeneration { 16 #[must_use] 17 pub const fn initial() -> Self { 18 Self(0) 19 } 20 21 #[must_use] 22 pub const fn from_value(value: u64) -> Self { 23 Self(value) 24 } 25 26 #[must_use] 27 pub const fn value(self) -> u64 { 28 self.0 29 } 30 31 #[must_use] 32 pub const fn next(self) -> Option<Self> { 33 match self.0.checked_add(1) { 34 Some(value) => Some(Self(value)), 35 None => None, 36 } 37 } 38 } 39 40 #[derive(Clone, Debug, Eq, PartialEq)] 41 pub struct ActiveSessionBinding { 42 identity: NostrIdentityReference, 43 signer_binding: SignerBinding, 44 generation: SessionGeneration, 45 } 46 47 impl ActiveSessionBinding { 48 /// Binds one foreground session to a ready local signer and generation. 49 /// 50 /// # Errors 51 /// 52 /// Returns a safe state error when identity and binding differ or when the 53 /// signer is unavailable. 54 pub fn new( 55 identity: NostrIdentityReference, 56 signer_binding: impl Into<SignerBinding>, 57 generation: SessionGeneration, 58 ) -> Result<Self, SafeError> { 59 let signer_binding = signer_binding.into(); 60 let local_keyring = signer_binding 61 .as_local_keyring() 62 .filter(|binding| binding.availability() == SignerAvailability::Available); 63 if local_keyring.map(|binding| binding.identity()) != Some(identity.public_key()) { 64 return Err(invalid_foreground_session()); 65 } 66 Ok(Self { 67 identity, 68 signer_binding, 69 generation, 70 }) 71 } 72 73 #[must_use] 74 pub const fn identity(&self) -> &NostrIdentityReference { 75 &self.identity 76 } 77 78 #[must_use] 79 pub const fn signer_binding(&self) -> SignerBinding { 80 self.signer_binding 81 } 82 83 #[must_use] 84 pub const fn generation(&self) -> SessionGeneration { 85 self.generation 86 } 87 } 88 89 const fn invalid_foreground_session() -> SafeError { 90 SafeError::new( 91 SafeErrorCode::InvalidApplicationState, 92 SafeMessage::new("The foreground session binding is invalid."), 93 ) 94 } 95 96 #[derive(Clone, Copy, Debug, Eq, PartialEq)] 97 pub struct TaskCorrelation { 98 request_id: RequestId, 99 identity: PublicKey, 100 binding: SignerBinding, 101 expected_revision: SnapshotRevision, 102 session_generation: SessionGeneration, 103 } 104 105 impl TaskCorrelation { 106 #[must_use] 107 pub const fn new( 108 request_id: RequestId, 109 identity: PublicKey, 110 binding: SignerBinding, 111 expected_revision: SnapshotRevision, 112 session_generation: SessionGeneration, 113 ) -> Self { 114 Self { 115 request_id, 116 identity, 117 binding, 118 expected_revision, 119 session_generation, 120 } 121 } 122 123 #[must_use] 124 pub const fn request_id(self) -> RequestId { 125 self.request_id 126 } 127 128 #[must_use] 129 pub const fn identity(self) -> PublicKey { 130 self.identity 131 } 132 133 #[must_use] 134 pub const fn binding(self) -> SignerBinding { 135 self.binding 136 } 137 138 #[must_use] 139 pub const fn expected_revision(self) -> SnapshotRevision { 140 self.expected_revision 141 } 142 143 #[must_use] 144 pub const fn session_generation(self) -> SessionGeneration { 145 self.session_generation 146 } 147 } 148 149 #[derive(Clone, Copy, Debug, Eq, PartialEq)] 150 pub enum RuntimeLifecycle { 151 Opening, 152 CompatibilityChecking, 153 AcquiringOwnership, 154 Migrating, 155 Recovering, 156 Ready, 157 Degraded(SafeError), 158 Blocked(SafeError), 159 ShuttingDown, 160 Closed, 161 Fatal(SafeError), 162 } 163 164 #[derive(Clone, Copy, Debug, Eq, PartialEq)] 165 pub enum RuntimeCommandClass { 166 Observe, 167 MutateLocalState, 168 UseCredential, 169 UseRelay, 170 RetryOpening, 171 Shutdown, 172 } 173 174 #[derive(Clone, Copy, Debug, Eq, PartialEq)] 175 pub struct LifecycleGate { 176 lifecycle: RuntimeLifecycle, 177 } 178 179 impl Default for LifecycleGate { 180 fn default() -> Self { 181 Self::opening() 182 } 183 } 184 185 impl LifecycleGate { 186 #[must_use] 187 pub const fn opening() -> Self { 188 Self { 189 lifecycle: RuntimeLifecycle::Opening, 190 } 191 } 192 193 #[must_use] 194 pub const fn lifecycle(self) -> RuntimeLifecycle { 195 self.lifecycle 196 } 197 198 #[must_use] 199 pub const fn allows(self, command: RuntimeCommandClass) -> bool { 200 match self.lifecycle { 201 RuntimeLifecycle::Opening 202 | RuntimeLifecycle::CompatibilityChecking 203 | RuntimeLifecycle::AcquiringOwnership 204 | RuntimeLifecycle::Migrating 205 | RuntimeLifecycle::Recovering => { 206 matches!( 207 command, 208 RuntimeCommandClass::Observe | RuntimeCommandClass::Shutdown 209 ) 210 } 211 RuntimeLifecycle::Ready => !matches!(command, RuntimeCommandClass::RetryOpening), 212 RuntimeLifecycle::Degraded(_) => !matches!( 213 command, 214 RuntimeCommandClass::UseRelay | RuntimeCommandClass::RetryOpening 215 ), 216 RuntimeLifecycle::Blocked(_) => matches!( 217 command, 218 RuntimeCommandClass::Observe 219 | RuntimeCommandClass::RetryOpening 220 | RuntimeCommandClass::Shutdown 221 ), 222 RuntimeLifecycle::ShuttingDown => matches!(command, RuntimeCommandClass::Observe), 223 RuntimeLifecycle::Closed | RuntimeLifecycle::Fatal(_) => false, 224 } 225 } 226 227 /// Advances the required open sequence to compatibility checking. 228 /// 229 /// # Errors 230 /// 231 /// Returns a safe lifecycle error when the stage is out of order. 232 pub fn begin_compatibility_check(&mut self) -> Result<(), SafeError> { 233 self.advance( 234 RuntimeLifecycle::Opening, 235 RuntimeLifecycle::CompatibilityChecking, 236 ) 237 } 238 239 /// Records compatibility acceptance and begins ownership acquisition. 240 /// 241 /// # Errors 242 /// 243 /// Returns a safe lifecycle error when the stage is out of order. 244 pub fn compatibility_accepted(&mut self) -> Result<(), SafeError> { 245 self.advance( 246 RuntimeLifecycle::CompatibilityChecking, 247 RuntimeLifecycle::AcquiringOwnership, 248 ) 249 } 250 251 /// Records exclusive ownership and begins migration. 252 /// 253 /// # Errors 254 /// 255 /// Returns a safe lifecycle error when the stage is out of order. 256 pub fn ownership_acquired(&mut self) -> Result<(), SafeError> { 257 self.advance( 258 RuntimeLifecycle::AcquiringOwnership, 259 RuntimeLifecycle::Migrating, 260 ) 261 } 262 263 /// Records migration completion and begins recovery. 264 /// 265 /// # Errors 266 /// 267 /// Returns a safe lifecycle error when the stage is out of order. 268 pub fn migration_complete(&mut self) -> Result<(), SafeError> { 269 self.advance(RuntimeLifecycle::Migrating, RuntimeLifecycle::Recovering) 270 } 271 272 /// Records recovery completion and admits normal commands. 273 /// 274 /// # Errors 275 /// 276 /// Returns a safe lifecycle error when the stage is out of order. 277 pub fn recovery_complete(&mut self) -> Result<(), SafeError> { 278 self.advance(RuntimeLifecycle::Recovering, RuntimeLifecycle::Ready) 279 } 280 281 pub fn block(&mut self, error: SafeError) { 282 self.lifecycle = RuntimeLifecycle::Blocked(error); 283 } 284 285 pub fn fail(&mut self, error: SafeError) { 286 self.lifecycle = RuntimeLifecycle::Fatal(error); 287 } 288 289 /// Moves a ready runtime into a nonfatal degraded state. 290 /// 291 /// # Errors 292 /// 293 /// Returns a safe lifecycle error when the runtime is not ready. 294 pub fn degrade(&mut self, error: SafeError) -> Result<(), SafeError> { 295 self.advance(RuntimeLifecycle::Ready, RuntimeLifecycle::Degraded(error)) 296 } 297 298 /// Restores local and relay command availability after degradation. 299 /// 300 /// # Errors 301 /// 302 /// Returns a safe lifecycle error when the runtime is not degraded. 303 pub fn restore_ready(&mut self) -> Result<(), SafeError> { 304 if !matches!(self.lifecycle, RuntimeLifecycle::Degraded(_)) { 305 return Err(invalid_lifecycle_transition()); 306 } 307 self.lifecycle = RuntimeLifecycle::Ready; 308 Ok(()) 309 } 310 311 /// Begins actor-owned shutdown. 312 /// 313 /// # Errors 314 /// 315 /// Returns a safe lifecycle error after shutdown or close has begun. 316 pub fn begin_shutdown(&mut self) -> Result<(), SafeError> { 317 if matches!( 318 self.lifecycle, 319 RuntimeLifecycle::ShuttingDown | RuntimeLifecycle::Closed 320 ) { 321 return Err(invalid_lifecycle_transition()); 322 } 323 self.lifecycle = RuntimeLifecycle::ShuttingDown; 324 Ok(()) 325 } 326 327 /// Completes actor-owned shutdown. 328 /// 329 /// # Errors 330 /// 331 /// Returns a safe lifecycle error unless shutdown already began. 332 pub fn finish_shutdown(&mut self) -> Result<(), SafeError> { 333 self.advance(RuntimeLifecycle::ShuttingDown, RuntimeLifecycle::Closed) 334 } 335 336 fn advance( 337 &mut self, 338 expected: RuntimeLifecycle, 339 next: RuntimeLifecycle, 340 ) -> Result<(), SafeError> { 341 if self.lifecycle != expected { 342 return Err(invalid_lifecycle_transition()); 343 } 344 self.lifecycle = next; 345 Ok(()) 346 } 347 } 348 349 const fn invalid_lifecycle_transition() -> SafeError { 350 SafeError::new( 351 harvestcircle_domain::SafeErrorCode::InvalidApplicationState, 352 harvestcircle_domain::SafeMessage::new("The runtime lifecycle transition is invalid."), 353 ) 354 } 355 356 #[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)] 357 pub struct RequestId(NonZeroU64); 358 359 impl RequestId { 360 #[must_use] 361 pub const fn new(value: u64) -> Option<Self> { 362 match NonZeroU64::new(value) { 363 Some(value) => Some(Self(value)), 364 None => None, 365 } 366 } 367 368 #[must_use] 369 pub const fn get(self) -> u64 { 370 self.0.get() 371 } 372 } 373 374 #[derive(Clone, Copy, Debug, Eq, PartialEq)] 375 pub struct CommandContext { 376 request_id: RequestId, 377 expected_revision: Option<SnapshotRevision>, 378 deadline: Instant, 379 } 380 381 impl CommandContext { 382 #[must_use] 383 pub const fn new( 384 request_id: RequestId, 385 expected_revision: Option<SnapshotRevision>, 386 deadline: Instant, 387 ) -> Self { 388 Self { 389 request_id, 390 expected_revision, 391 deadline, 392 } 393 } 394 395 #[must_use] 396 pub const fn request_id(self) -> RequestId { 397 self.request_id 398 } 399 400 #[must_use] 401 pub const fn expected_revision(self) -> Option<SnapshotRevision> { 402 self.expected_revision 403 } 404 405 #[must_use] 406 pub const fn deadline(self) -> Instant { 407 self.deadline 408 } 409 410 #[must_use] 411 pub fn is_expired(self, now: Instant) -> bool { 412 now >= self.deadline 413 } 414 } 415 416 #[derive(Clone, Copy, Debug, Eq, PartialEq)] 417 pub enum CommandRejection { 418 MailboxSaturated, 419 } 420 421 #[derive(Clone, Debug, Eq, PartialEq)] 422 pub enum CommandResult<T> { 423 Completed(T), 424 Rejected(CommandRejection), 425 Conflicted { current_revision: SnapshotRevision }, 426 TimedOut, 427 Closed, 428 Failed(SafeError), 429 } 430 431 #[derive(Clone, Debug, Eq, PartialEq)] 432 pub struct CommandReceipt<T> { 433 request_id: RequestId, 434 result: CommandResult<T>, 435 } 436 437 impl<T> CommandReceipt<T> { 438 #[must_use] 439 pub const fn new(request_id: RequestId, result: CommandResult<T>) -> Self { 440 Self { request_id, result } 441 } 442 443 #[must_use] 444 pub const fn request_id(&self) -> RequestId { 445 self.request_id 446 } 447 448 #[must_use] 449 pub const fn result(&self) -> &CommandResult<T> { 450 &self.result 451 } 452 453 #[must_use] 454 pub fn into_result(self) -> CommandResult<T> { 455 self.result 456 } 457 } 458 459 pub struct CommandTicket<T> { 460 request_id: RequestId, 461 receiver: oneshot::Receiver<CommandReceipt<T>>, 462 } 463 464 impl<T> CommandTicket<T> { 465 #[must_use] 466 pub const fn request_id(&self) -> RequestId { 467 self.request_id 468 } 469 470 pub async fn receipt(self) -> CommandReceipt<T> { 471 self.receiver 472 .await 473 .unwrap_or_else(|_| CommandReceipt::new(self.request_id, CommandResult::Closed)) 474 } 475 } 476 477 pub enum CommandSubmission<T> { 478 Accepted(CommandTicket<T>), 479 Rejected(CommandReceipt<T>), 480 } 481 482 impl<T> CommandSubmission<T> { 483 #[must_use] 484 pub const fn request_id(&self) -> RequestId { 485 match self { 486 Self::Accepted(ticket) => ticket.request_id(), 487 Self::Rejected(receipt) => receipt.request_id(), 488 } 489 } 490 } 491 492 pub struct CommandEnvelope<C, R> { 493 context: CommandContext, 494 command: C, 495 reply: oneshot::Sender<CommandReceipt<R>>, 496 } 497 498 impl<C, R> CommandEnvelope<C, R> { 499 #[must_use] 500 pub const fn context(&self) -> CommandContext { 501 self.context 502 } 503 504 #[must_use] 505 pub const fn command(&self) -> &C { 506 &self.command 507 } 508 509 #[must_use] 510 pub fn into_parts(self) -> (CommandContext, C, oneshot::Sender<CommandReceipt<R>>) { 511 (self.context, self.command, self.reply) 512 } 513 } 514 515 pub struct ActorMailbox<C, R> { 516 sender: mpsc::Sender<CommandEnvelope<C, R>>, 517 } 518 519 impl<C, R> Clone for ActorMailbox<C, R> { 520 fn clone(&self) -> Self { 521 Self { 522 sender: self.sender.clone(), 523 } 524 } 525 } 526 527 impl<C, R> ActorMailbox<C, R> { 528 #[must_use] 529 pub fn bounded(capacity: NonZeroUsize) -> (Self, mpsc::Receiver<CommandEnvelope<C, R>>) { 530 let (sender, receiver) = mpsc::channel(capacity.get()); 531 (Self { sender }, receiver) 532 } 533 534 #[must_use] 535 pub fn available_capacity(&self) -> usize { 536 self.sender.capacity() 537 } 538 539 #[must_use] 540 pub fn submit(&self, context: CommandContext, command: C) -> CommandSubmission<R> { 541 let request_id = context.request_id(); 542 if context.is_expired(Instant::now()) { 543 return CommandSubmission::Rejected(CommandReceipt::new( 544 request_id, 545 CommandResult::TimedOut, 546 )); 547 } 548 let (reply, receiver) = oneshot::channel(); 549 let envelope = CommandEnvelope { 550 context, 551 command, 552 reply, 553 }; 554 match self.sender.try_send(envelope) { 555 Ok(()) => CommandSubmission::Accepted(CommandTicket { 556 request_id, 557 receiver, 558 }), 559 Err(mpsc::error::TrySendError::Full(_)) => { 560 CommandSubmission::Rejected(CommandReceipt::new( 561 request_id, 562 CommandResult::Rejected(CommandRejection::MailboxSaturated), 563 )) 564 } 565 Err(mpsc::error::TrySendError::Closed(_)) => { 566 CommandSubmission::Rejected(CommandReceipt::new(request_id, CommandResult::Closed)) 567 } 568 } 569 } 570 } 571 572 #[cfg(test)] 573 mod tests { 574 use std::num::NonZeroUsize; 575 use std::time::{Duration, Instant}; 576 577 use harvestcircle_domain::{LocalKeyringBinding, NostrIdentityReference, SignerAvailability}; 578 579 use crate::{ 580 ActiveSessionBinding, ActorMailbox, CommandContext, CommandReceipt, CommandRejection, 581 CommandResult, CommandSubmission, LifecycleGate, RequestId, RuntimeCommandClass, 582 RuntimeLifecycle, SessionGeneration, 583 }; 584 585 fn context(id: u64) -> CommandContext { 586 CommandContext::new( 587 RequestId::new(id).expect("nonzero request"), 588 None, 589 Instant::now() + Duration::from_secs(1), 590 ) 591 } 592 593 #[test] 594 fn foreground_session_requires_matching_available_binding_and_generation() { 595 let public_key = crate::test_support::valid_test_public_key(3).expect("valid public key"); 596 let identity = NostrIdentityReference::derive(public_key).expect("identity"); 597 let generation = SessionGeneration::from_value(4); 598 let session = ActiveSessionBinding::new( 599 identity.clone(), 600 LocalKeyringBinding::new(public_key, SignerAvailability::Available), 601 generation, 602 ) 603 .expect("session"); 604 assert_eq!(session.identity(), &identity); 605 assert_eq!(session.signer_binding().identity(), public_key); 606 assert!(session.signer_binding().as_local_keyring().is_some()); 607 assert_eq!(session.generation(), generation); 608 609 let correlation = super::TaskCorrelation::new( 610 RequestId::new(8).expect("request"), 611 public_key, 612 session.signer_binding(), 613 crate::SnapshotRevision::from_value(9), 614 generation, 615 ); 616 assert_eq!(correlation.request_id().get(), 8); 617 assert_eq!(correlation.identity(), public_key); 618 assert_eq!(correlation.binding(), session.signer_binding()); 619 assert_eq!(correlation.expected_revision().value(), 9); 620 assert_eq!(correlation.session_generation(), generation); 621 622 assert!( 623 ActiveSessionBinding::new( 624 identity.clone(), 625 LocalKeyringBinding::new( 626 harvestcircle_domain::PublicKey::from_hex( 627 "e0266e3cfb0d2886f91c73f5f868f3b98273713e5fcd97c081663f5518a4b3af", 628 ) 629 .expect("different valid public key"), 630 SignerAvailability::Available, 631 ), 632 generation, 633 ) 634 .is_err() 635 ); 636 assert!( 637 ActiveSessionBinding::new( 638 identity, 639 LocalKeyringBinding::new(public_key, SignerAvailability::CredentialMissing), 640 generation, 641 ) 642 .is_err() 643 ); 644 } 645 646 #[tokio::test] 647 async fn bounded_mailbox_accepts_one_and_rejects_saturation() { 648 let (mailbox, mut receiver) = 649 ActorMailbox::<u8, u8>::bounded(NonZeroUsize::new(1).expect("capacity")); 650 let CommandSubmission::Accepted(ticket) = mailbox.submit(context(1), 7) else { 651 panic!("first command must be accepted"); 652 }; 653 let CommandSubmission::Rejected(rejected) = mailbox.submit(context(2), 8) else { 654 panic!("second command must be rejected"); 655 }; 656 assert_eq!( 657 rejected.into_result(), 658 CommandResult::Rejected(CommandRejection::MailboxSaturated) 659 ); 660 661 let envelope = receiver.recv().await.expect("command"); 662 assert_eq!(envelope.context().request_id().get(), 1); 663 assert_eq!(*envelope.command(), 7); 664 let (context, command, reply) = envelope.into_parts(); 665 reply 666 .send(CommandReceipt::new( 667 context.request_id(), 668 CommandResult::Completed(command + 1), 669 )) 670 .expect("ticket remains open"); 671 assert_eq!( 672 ticket.receipt().await.into_result(), 673 CommandResult::Completed(8) 674 ); 675 } 676 677 #[test] 678 fn expired_and_closed_mailboxes_reject_without_enqueuing() { 679 let (mailbox, receiver) = 680 ActorMailbox::<u8, u8>::bounded(NonZeroUsize::new(1).expect("capacity")); 681 let expired = 682 CommandContext::new(RequestId::new(1).expect("request"), None, Instant::now()); 683 let CommandSubmission::Rejected(receipt) = mailbox.submit(expired, 1) else { 684 panic!("expired command must be rejected"); 685 }; 686 assert_eq!(receipt.into_result(), CommandResult::TimedOut); 687 688 drop(receiver); 689 let CommandSubmission::Rejected(receipt) = mailbox.submit(context(2), 2) else { 690 panic!("closed mailbox must be rejected"); 691 }; 692 assert_eq!(receipt.into_result(), CommandResult::Closed); 693 } 694 695 #[tokio::test] 696 async fn dropped_actor_reply_becomes_closed_receipt() { 697 let (mailbox, mut receiver) = 698 ActorMailbox::<u8, u8>::bounded(NonZeroUsize::new(1).expect("capacity")); 699 let CommandSubmission::Accepted(ticket) = mailbox.submit(context(1), 1) else { 700 panic!("command must be accepted"); 701 }; 702 drop(receiver.recv().await.expect("command")); 703 704 assert_eq!(ticket.receipt().await.into_result(), CommandResult::Closed); 705 } 706 707 #[test] 708 fn opening_sequence_gates_mutation_until_recovery_completes() { 709 let mut lifecycle = LifecycleGate::opening(); 710 for expected in [ 711 RuntimeLifecycle::Opening, 712 RuntimeLifecycle::CompatibilityChecking, 713 RuntimeLifecycle::AcquiringOwnership, 714 RuntimeLifecycle::Migrating, 715 RuntimeLifecycle::Recovering, 716 ] { 717 assert_eq!(lifecycle.lifecycle(), expected); 718 assert!(lifecycle.allows(RuntimeCommandClass::Observe)); 719 assert!(lifecycle.allows(RuntimeCommandClass::Shutdown)); 720 assert!(!lifecycle.allows(RuntimeCommandClass::MutateLocalState)); 721 match expected { 722 RuntimeLifecycle::Opening => { 723 lifecycle 724 .begin_compatibility_check() 725 .expect("compatibility"); 726 } 727 RuntimeLifecycle::CompatibilityChecking => { 728 lifecycle 729 .compatibility_accepted() 730 .expect("compatibility accepted"); 731 } 732 RuntimeLifecycle::AcquiringOwnership => { 733 lifecycle.ownership_acquired().expect("ownership"); 734 } 735 RuntimeLifecycle::Migrating => { 736 lifecycle.migration_complete().expect("migration"); 737 } 738 RuntimeLifecycle::Recovering => { 739 lifecycle.recovery_complete().expect("recovery"); 740 } 741 _ => unreachable!("opening states only"), 742 } 743 } 744 assert_eq!(lifecycle.lifecycle(), RuntimeLifecycle::Ready); 745 assert!(lifecycle.allows(RuntimeCommandClass::MutateLocalState)); 746 assert!(lifecycle.allows(RuntimeCommandClass::UseCredential)); 747 assert!(lifecycle.allows(RuntimeCommandClass::UseRelay)); 748 } 749 750 #[test] 751 fn blocked_degraded_fatal_and_closed_states_fail_safe() { 752 let problem = harvestcircle_domain::SafeError::new( 753 harvestcircle_domain::SafeErrorCode::StorageUnavailable, 754 harvestcircle_domain::SafeMessage::new("The runtime is unavailable."), 755 ); 756 let mut blocked = LifecycleGate::opening(); 757 blocked.block(problem); 758 assert!(blocked.allows(RuntimeCommandClass::RetryOpening)); 759 assert!(!blocked.allows(RuntimeCommandClass::MutateLocalState)); 760 761 let mut degraded = LifecycleGate::opening(); 762 degraded.begin_compatibility_check().expect("compatibility"); 763 degraded.compatibility_accepted().expect("accepted"); 764 degraded.ownership_acquired().expect("ownership"); 765 degraded.migration_complete().expect("migration"); 766 degraded.recovery_complete().expect("recovery"); 767 degraded.degrade(problem).expect("degraded"); 768 assert!(degraded.allows(RuntimeCommandClass::MutateLocalState)); 769 assert!(!degraded.allows(RuntimeCommandClass::UseRelay)); 770 degraded.restore_ready().expect("restored"); 771 assert!(LifecycleGate::opening().restore_ready().is_err()); 772 773 let mut fatal = LifecycleGate::opening(); 774 fatal.fail(problem); 775 assert!(!fatal.allows(RuntimeCommandClass::Observe)); 776 fatal.begin_shutdown().expect("fatal can close"); 777 fatal.finish_shutdown().expect("closed"); 778 assert_eq!(fatal.lifecycle(), RuntimeLifecycle::Closed); 779 assert!(!fatal.allows(RuntimeCommandClass::Shutdown)); 780 } 781 782 #[test] 783 fn opening_stages_reject_out_of_order_and_repeated_transitions() { 784 let mut lifecycle = LifecycleGate::opening(); 785 assert!(lifecycle.migration_complete().is_err()); 786 lifecycle.begin_compatibility_check().expect("first stage"); 787 assert!(lifecycle.begin_compatibility_check().is_err()); 788 assert!(lifecycle.recovery_complete().is_err()); 789 } 790 }