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 }