lib

Core libraries for Radroots
git clone https://radroots.dev/git/lib.git
Log | Files | Refs | README

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 }