capacity.rs (4835B)
1 use super::*; 2 use radroots_storage::authored_atomic::AuthoredAtomicCommand; 3 4 #[derive(Clone, Copy)] 5 pub(super) enum Phase { 6 Prepare, 7 Claim, 8 Signed, 9 DeliveryFact, 10 } 11 12 pub(super) struct Fault { 13 pub(super) phase: Phase, 14 pub(super) after: bool, 15 } 16 17 pub(super) fn inject( 18 storage: &FaultStorage, 19 command: &AuthoredAtomicCommand, 20 after: bool, 21 ) -> Result<(), radroots_storage::Error> { 22 let mut fault = storage.capacity.lock().unwrap(); 23 if fault.as_ref().is_some_and(|fault| { 24 fault.after == after 25 && matches!( 26 (fault.phase, command), 27 (Phase::Prepare, AuthoredAtomicCommand::Prepare(_)) 28 | (Phase::Claim, AuthoredAtomicCommand::Claim(_)) 29 | (Phase::Signed, AuthoredAtomicCommand::RecordSigned(_)) 30 | ( 31 Phase::DeliveryFact, 32 AuthoredAtomicCommand::RecordDelivery(_) 33 ) 34 ) 35 }) { 36 *fault = None; 37 Err(radroots_storage::Error::SpaceInsufficient) 38 } else { 39 Ok(()) 40 } 41 } 42 43 fn setup(byte: u8) -> (Engine, Arc<FaultStorage>, Arc<MockSigner>, PushRequest) { 44 let storage = Arc::new(FaultStorage::new(byte)); 45 let signer = Arc::new(MockSigner::new(SignBehavior::Success { 46 completed_at_unix_ms: 1_800_000_200_500, 47 })); 48 let engine = fault_engine(storage.clone(), signer.clone(), Arc::new(MockSink)); 49 ( 50 engine, 51 storage, 52 signer, 53 request(byte, "wss://capacity.example"), 54 ) 55 } 56 57 #[test] 58 fn capacity_prepare_reports_failure_and_reconciles_the_original_operation() { 59 for after in [false, true] { 60 let (engine, storage, signer, push) = setup(181); 61 *storage.capacity.lock().unwrap() = Some(Fault { 62 phase: Phase::Prepare, 63 after, 64 }); 65 assert_eq!( 66 block_on(engine.prepare_push(push.clone())), 67 Err(Error::StorageSpaceInsufficient) 68 ); 69 let observed = block_on(engine.push_status(push.operation_id())).unwrap(); 70 assert_eq!(observed.is_some(), after); 71 let original = block_on(engine.prepare_push(push.clone())).unwrap(); 72 let replay = block_on(engine.prepare_push(push)).unwrap(); 73 assert_eq!(original.operation(), replay.operation()); 74 assert_eq!(original.artifact(), replay.artifact()); 75 assert_eq!(signer.calls.load(Ordering::Relaxed), 0); 76 } 77 } 78 79 #[test] 80 fn capacity_after_signed_receipt_preserves_exact_bytes_without_resigning() { 81 let (engine, storage, signer, push) = setup(182); 82 *storage.capacity.lock().unwrap() = Some(Fault { 83 phase: Phase::Signed, 84 after: true, 85 }); 86 assert_eq!( 87 block_on(engine.sign_prepared(push.clone())), 88 Err(Error::StorageSpaceInsufficient) 89 ); 90 let original = block_on(engine.push_status(push.operation_id())) 91 .unwrap() 92 .unwrap(); 93 let signed = original 94 .artifact() 95 .signed() 96 .expect("durable original signed bytes"); 97 let replay = block_on(engine.sign_prepared(push)).unwrap(); 98 assert_eq!(replay.artifact().signed().unwrap(), signed); 99 assert_eq!(signer.calls.load(Ordering::Relaxed), 1); 100 } 101 102 #[test] 103 fn capacity_claim_failure_never_calls_the_signer_or_discards_prepared_work() { 104 for after in [false, true] { 105 let (engine, storage, signer, push) = setup(183); 106 block_on(engine.prepare_push(push.clone())).unwrap(); 107 *storage.capacity.lock().unwrap() = Some(Fault { 108 phase: Phase::Claim, 109 after, 110 }); 111 assert_eq!( 112 block_on(engine.sign_prepared(push.clone())), 113 Err(Error::StorageSpaceInsufficient) 114 ); 115 let status = block_on(engine.push_status(push.operation_id())) 116 .unwrap() 117 .unwrap(); 118 assert_eq!(status.artifact().signing_claim().is_some(), after); 119 assert!(status.artifact().signed().is_none()); 120 assert_eq!(signer.calls.load(Ordering::Relaxed), 0); 121 } 122 } 123 124 #[test] 125 fn capacity_local_admission_retains_signed_bytes_and_reports_no_success() { 126 let (engine, storage, signer, push) = setup(184); 127 block_on(engine.sign_prepared(push.clone())).unwrap(); 128 let original = block_on(engine.push_status(push.operation_id())) 129 .unwrap() 130 .unwrap(); 131 storage.fail_admission_with(3); 132 assert_eq!( 133 block_on(engine.admit_signed(push.operation_id())), 134 Err(Error::StorageSpaceInsufficient) 135 ); 136 let status = block_on(engine.push_status(push.operation_id())) 137 .unwrap() 138 .unwrap(); 139 assert_eq!(status.artifact().signed(), original.artifact().signed()); 140 assert!(!status.artifact().admission_state().is_admitted()); 141 assert_eq!(signer.calls.load(Ordering::Relaxed), 1); 142 }