lib

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

signing_evidence.rs (14520B)


      1 use super::*;
      2 use futures::channel::oneshot;
      3 use radroots_signing::{AuthoredSignEvidence, SigningIntentId, SigningOperationId};
      4 use radroots_storage::{
      5     authored::{AdmissionState, FailureClass, SigningState, WorkFailure, WorkPhase},
      6     authored_atomic::{ApplyWorkFailure, AuthoredAtomicCommand, AuthoredWorkTarget, WorkFence},
      7 };
      8 
      9 const NOW: u64 = 1_800_000_200_000;
     10 type Pending = (
     11     SignRequest,
     12     oneshot::Sender<Result<AuthoredSignEvidence, SigningError>>,
     13 );
     14 
     15 #[derive(Default)]
     16 struct HeldSigner {
     17     pending: Mutex<VecDeque<Pending>>,
     18     evidence_calls: AtomicUsize,
     19     legacy_calls: AtomicUsize,
     20 }
     21 
     22 impl HeldSigner {
     23     fn take(&self) -> Pending {
     24         self.pending
     25             .lock()
     26             .unwrap()
     27             .pop_front()
     28             .expect("started request")
     29     }
     30 }
     31 
     32 impl Signer for HeldSigner {
     33     fn status(
     34         &self,
     35     ) -> radroots_signing::signer::BoxFuture<'_, Result<SignerStatus, SigningError>> {
     36         Box::pin(async {
     37             Ok(SignerStatus::new(
     38                 SignerAvailability::Ready,
     39                 vec![SignerCapability::new(
     40                     SignerKind::Remote,
     41                     ReplayCapability::ExactReplayByRequestId,
     42                     CancellationSupport::BeforeAndAfterPublication,
     43                     false,
     44                     false,
     45                 )],
     46                 None,
     47             ))
     48         })
     49     }
     50 
     51     fn sign(
     52         &self,
     53         _: SignRequest,
     54     ) -> radroots_signing::signer::BoxFuture<'_, Result<SignReceipt, SigningError>> {
     55         self.legacy_calls.fetch_add(1, Ordering::Relaxed);
     56         Box::pin(async { Err(SigningError::new(SigningErrorKind::InternalError)) })
     57     }
     58 
     59     fn sign_authored_evidence(
     60         &self,
     61         request: SignRequest,
     62     ) -> radroots_signing::signer::BoxFuture<'_, Result<AuthoredSignEvidence, SigningError>> {
     63         Box::pin(async move {
     64             self.evidence_calls.fetch_add(1, Ordering::Relaxed);
     65             let (sender, receiver) = oneshot::channel();
     66             self.pending.lock().unwrap().push_back((request, sender));
     67             receiver
     68                 .await
     69                 .map_err(|_| SigningError::new(SigningErrorKind::SignerUnavailable))?
     70         })
     71     }
     72 }
     73 
     74 fn engine(
     75     storage: Arc<dyn SyncStorage>,
     76     clock: Arc<dyn Clock>,
     77     signer: Option<Arc<dyn Signer>>,
     78 ) -> Engine {
     79     let builder = Engine::builder(
     80         storage,
     81         clock,
     82         Arc::new(TestIds(AtomicU64::new(100))),
     83         DeadlinePolicy::new(10_000, 10_000, 10_000).unwrap(),
     84     )
     85     .sink(Arc::new(MockSink));
     86     if let Some(signer) = signer {
     87         builder.signer(signer).build().unwrap()
     88     } else {
     89         builder.build().unwrap()
     90     }
     91 }
     92 
     93 fn setup() -> (Engine, Arc<MemoryStorage>, Arc<TestClock>, Arc<HeldSigner>) {
     94     let storage = Arc::new(MemoryStorage::new(SourceGeneration::new([17; 32]).unwrap()));
     95     let clock = Arc::new(TestClock(AtomicU64::new(NOW)));
     96     let signer = Arc::new(HeldSigner::default());
     97     (
     98         engine(storage.clone(), clock.clone(), Some(signer.clone())),
     99         storage,
    100         clock,
    101         signer,
    102     )
    103 }
    104 
    105 fn complete((request, sender): Pending, at: u64) -> SignedEvent {
    106     let event = signed_event(&request);
    107     let evidence = AuthoredSignEvidence::from_signed_event(&request, event.clone(), at).unwrap();
    108     sender.send(Ok(evidence)).unwrap();
    109     event
    110 }
    111 
    112 #[test]
    113 fn missing_observation_clock_preserves_uncertain_non_replayable_attempt() {
    114     struct FaultClock(AtomicUsize);
    115     impl Clock for FaultClock {
    116         fn now_unix_ms(&self) -> Result<u64, Error> {
    117             match self.0.fetch_add(1, Ordering::Relaxed) {
    118                 2 => Err(Error::ClockUnavailable),
    119                 0 | 1 => Ok(NOW),
    120                 _ => Ok(NOW + 11_000),
    121             }
    122         }
    123     }
    124     let storage = Arc::new(MemoryStorage::new(SourceGeneration::new([37; 32]).unwrap()));
    125     let signer = Arc::new(MockSigner::with_replay(
    126         SignBehavior::Success {
    127             completed_at_unix_ms: NOW,
    128         },
    129         ReplayCapability::NonReplayable,
    130     ));
    131     let engine = engine(
    132         storage,
    133         Arc::new(FaultClock(AtomicUsize::new(0))),
    134         Some(signer.clone()),
    135     );
    136     let push = request(37, "wss://relay.example");
    137     assert_eq!(
    138         block_on(engine.sign_prepared(push.clone())),
    139         Err(Error::ClockUnavailable)
    140     );
    141     let status = block_on(engine.push_status(push.operation_id()))
    142         .unwrap()
    143         .unwrap();
    144     assert!(status.artifact().signing_claim().is_some());
    145     assert!(status.artifact().signed().is_none());
    146     assert!(status.artifact().last_failure().is_none());
    147     assert_eq!(
    148         block_on(engine.sign_prepared(push)),
    149         Err(Error::SigningIndeterminate)
    150     );
    151     assert_eq!(signer.calls.load(Ordering::Relaxed), 1);
    152 }
    153 
    154 #[test]
    155 fn late_evidence_survives_durable_cancel_and_cannot_schedule_more_work() {
    156     let (engine, storage, clock, signer) = setup();
    157     let push = request(31, "wss://relay.example");
    158     let mut future = Box::pin(engine.sign_prepared(push.clone())).fuse();
    159     assert!(
    160         future
    161             .poll_unpin(&mut std::task::Context::from_waker(noop_waker_ref()))
    162             .is_pending()
    163     );
    164     let pending = signer.take();
    165     block_on(engine.cancel_push(push.operation_id())).unwrap();
    166     clock.0.store(NOW + 11_000, Ordering::Release);
    167     let event = complete(pending, NOW + 11_000);
    168     assert_eq!(block_on(future), Err(Error::SigningCancelled));
    169     let status = block_on(engine.push_status(push.operation_id()))
    170         .unwrap()
    171         .unwrap();
    172     assert_eq!(status.artifact().signing_state(), SigningState::Cancelled);
    173     assert_eq!(status.artifact().signed().unwrap().event(), &event);
    174     assert_eq!(
    175         status.delivery_plan().state(),
    176         AuthoredDeliveryState::Cancelled
    177     );
    178     assert!(block_on(engine.admit_signed(push.operation_id())).is_err());
    179     assert_eq!(
    180         block_on(engine.sign_prepared(push)),
    181         Err(Error::SigningCancelled)
    182     );
    183     assert_eq!(signer.evidence_calls.load(Ordering::Relaxed), 1);
    184     assert_eq!(signer.legacy_calls.load(Ordering::Relaxed), 0);
    185     assert!(
    186         block_on(storage.query_visible(EventQuery::all(EventQueryBounds::first(10).unwrap())))
    187             .unwrap()
    188             .items()
    189             .is_empty()
    190     );
    191 }
    192 
    193 #[test]
    194 fn superseded_attempts_retain_first_bytes_and_replay_without_a_signer() {
    195     let (engine, storage, clock, signer) = setup();
    196     let push = request(32, "wss://relay.example");
    197     let mut first = Box::pin(engine.sign_prepared(push.clone())).fuse();
    198     let mut context = std::task::Context::from_waker(noop_waker_ref());
    199     assert!(first.poll_unpin(&mut context).is_pending());
    200     let first_result = signer.take();
    201     clock.0.store(NOW + 11_000, Ordering::Release);
    202     let mut second = Box::pin(engine.sign_prepared(push.clone())).fuse();
    203     assert!(second.poll_unpin(&mut context).is_pending());
    204     let second_result = signer.take();
    205     assert_eq!(
    206         first_result.0.signer_request_id(),
    207         second_result.0.signer_request_id()
    208     );
    209     let event = complete(first_result, NOW + 11_000);
    210     assert_eq!(block_on(first), Err(Error::SignerDeadlineExceeded));
    211     let retained = block_on(engine.push_status(push.operation_id()))
    212         .unwrap()
    213         .unwrap();
    214     assert_eq!(complete(second_result, NOW + 11_000), event);
    215     let result = block_on(second).unwrap();
    216     assert_eq!(result.artifact(), retained.artifact());
    217     assert_eq!(result.artifact().signed().unwrap().event(), &event);
    218     let recovered = self::engine(storage, clock, None);
    219     let replay = block_on(recovered.sign_prepared(push)).unwrap();
    220     assert!(replay.is_replay());
    221     assert_eq!(replay.artifact(), retained.artifact());
    222     assert_eq!(signer.evidence_calls.load(Ordering::Relaxed), 2);
    223     assert_eq!(signer.legacy_calls.load(Ordering::Relaxed), 0);
    224 }
    225 
    226 #[test]
    227 fn stale_signer_failure_cannot_overwrite_another_workers_signed_fact() {
    228     let (engine, _, clock, signer) = setup();
    229     let push = request(33, "wss://relay.example");
    230     let mut first = Box::pin(engine.sign_prepared(push.clone())).fuse();
    231     let mut context = std::task::Context::from_waker(noop_waker_ref());
    232     assert!(first.poll_unpin(&mut context).is_pending());
    233     let first_result = signer.take();
    234     clock.0.store(NOW + 11_000, Ordering::Release);
    235     let mut second = Box::pin(engine.sign_prepared(push.clone())).fuse();
    236     assert!(second.poll_unpin(&mut context).is_pending());
    237     complete(signer.take(), NOW + 11_000);
    238     let signed = block_on(second).unwrap();
    239     first_result
    240         .1
    241         .send(Err(SigningError::new(SigningErrorKind::SignerRejected)))
    242         .unwrap();
    243     assert_eq!(block_on(first), Err(Error::SignerFailed));
    244     let status = block_on(engine.push_status(push.operation_id()))
    245         .unwrap()
    246         .unwrap();
    247     assert_eq!(status.artifact(), signed.artifact());
    248     assert!(status.artifact().last_failure().is_none());
    249 }
    250 
    251 #[test]
    252 fn valid_signature_under_another_operation_or_artifact_is_rejected() {
    253     for change_operation in [false, true] {
    254         let (engine, _, _, signer) = setup();
    255         let push = request(34, "wss://relay.example");
    256         let mut future = Box::pin(engine.sign_prepared(push.clone())).fuse();
    257         assert!(
    258             future
    259                 .poll_unpin(&mut std::task::Context::from_waker(noop_waker_ref()))
    260                 .is_pending()
    261         );
    262         let (original, sender) = signer.take();
    263         let intent = if change_operation {
    264             SigningIntentId::new(
    265                 SigningOperationId::new([99; 16]).unwrap(),
    266                 original.intent_id().artifact_id(),
    267             )
    268         } else {
    269             SigningIntentId::new(
    270                 original.intent_id().operation_id(),
    271                 radroots_signing::AuthoredArtifactId::new([99; 16]).unwrap(),
    272             )
    273         };
    274         let other = SignRequest::new(
    275             original.operation_kind(),
    276             intent,
    277             original.actor().clone(),
    278             original.authored_plan().unwrap().clone(),
    279             original.policy(),
    280         )
    281         .unwrap();
    282         let evidence =
    283             AuthoredSignEvidence::from_signed_event(&other, signed_event(&other), NOW).unwrap();
    284         sender.send(Ok(evidence)).unwrap();
    285         assert_eq!(block_on(future), Err(Error::SignerFailed));
    286         let status = block_on(engine.push_status(push.operation_id()))
    287             .unwrap()
    288             .unwrap();
    289         assert!(status.artifact().signed().is_none());
    290         assert!(status.delivery_plan().request().is_none());
    291         assert!(status.delivery_plan().attempts().is_empty());
    292     }
    293 }
    294 
    295 #[test]
    296 fn cancelled_indeterminate_operation_keeps_late_facts_without_admission() {
    297     let (engine, storage, clock, signer) = setup();
    298     let push = request(35, "wss://relay.example");
    299     let mut future = Box::pin(engine.sign_prepared(push.clone())).fuse();
    300     assert!(
    301         future
    302             .poll_unpin(&mut std::task::Context::from_waker(noop_waker_ref()))
    303             .is_pending()
    304     );
    305     let pending = signer.take();
    306     let status = block_on(engine.push_status(push.operation_id()))
    307         .unwrap()
    308         .unwrap();
    309     let claim = status.artifact().signing_claim().unwrap();
    310     block_on(
    311         storage.execute_authored(AuthoredAtomicCommand::ApplyFailure(
    312             ApplyWorkFailure::new(
    313                 AuthoredWorkTarget::Artifact(status.artifact().artifact_id()),
    314                 WorkFence::new(*claim.token(), claim.generation(), claim.row_revision()).unwrap(),
    315                 WorkFailure::new(
    316                     "signing_unknown",
    317                     WorkPhase::Signing,
    318                     FailureClass::Indeterminate,
    319                     None,
    320                     None,
    321                 )
    322                 .unwrap(),
    323                 None,
    324                 NOW + 1,
    325             )
    326             .unwrap(),
    327         )),
    328     )
    329     .unwrap();
    330     clock.0.store(NOW + 2, Ordering::Release);
    331     let stopped = block_on(engine.cancel_push(push.operation_id())).unwrap();
    332     assert_eq!(
    333         stopped.status().artifact().signing_state(),
    334         SigningState::Indeterminate
    335     );
    336     let event = complete(pending, NOW + 2);
    337     assert_eq!(block_on(future), Err(Error::SigningCancelled));
    338     let status = block_on(engine.push_status(push.operation_id()))
    339         .unwrap()
    340         .unwrap();
    341     assert_eq!(status.artifact().signed().unwrap().event(), &event);
    342     assert_eq!(
    343         status.artifact().admission_state(),
    344         AdmissionState::Cancelled
    345     );
    346     assert_eq!(
    347         block_on(engine.admit_signed(push.operation_id())),
    348         Err(Error::AdmissionFailed)
    349     );
    350     assert!(status.delivery_plan().attempts().is_empty());
    351 }
    352 
    353 #[tokio::test]
    354 async fn late_sqlite_fact_reopens_without_requiring_missing_credentials() {
    355     let directory = tempfile::tempdir().unwrap();
    356     let paths = Paths::from_directory(directory.path()).unwrap();
    357     let clock = Arc::new(TestClock(AtomicU64::new(NOW)));
    358     let store = Arc::new(
    359         SqliteStorage::open(
    360             OpenOptions::new(paths.clone(), OpenMode::Create)
    361                 .with_source_generation(SourceGeneration::new([36; 32]).unwrap(), 1)
    362                 .unwrap(),
    363         )
    364         .await
    365         .unwrap(),
    366     );
    367     let signer = Arc::new(BoundaryViolatingSigner {
    368         violation: BoundaryViolation::CompletesAfterDeadline,
    369         clock: clock.clone(),
    370     });
    371     let original = engine(store.clone(), clock.clone(), Some(signer));
    372     let push = request(36, "wss://relay.example");
    373     assert_eq!(
    374         original.sign_prepared(push.clone()).await,
    375         Err(Error::SignerDeadlineExceeded)
    376     );
    377     let before = original
    378         .push_status(push.operation_id())
    379         .await
    380         .unwrap()
    381         .unwrap();
    382     assert!(before.artifact().signed().is_some());
    383     drop(original);
    384     store.close().await.unwrap();
    385     let store = Arc::new(
    386         SqliteStorage::open(OpenOptions::new(paths, OpenMode::ReadWriteExisting))
    387             .await
    388             .unwrap(),
    389     );
    390     let unavailable = Arc::new(MockSigner::new(SignBehavior::Error(
    391         SigningErrorKind::SignerUnavailable,
    392     )));
    393     let recovered = engine(store.clone(), clock, Some(unavailable.clone()));
    394     let replay = recovered.sign_prepared(push).await.unwrap();
    395     assert!(replay.is_replay());
    396     assert_eq!(replay.artifact(), before.artifact());
    397     assert_eq!(unavailable.calls.load(Ordering::Relaxed), 0);
    398     drop(recovered);
    399     store.close().await.unwrap();
    400 }