authored_signed_facts.rs (16451B)
1 use std::num::NonZeroU64; 2 3 use futures_executor::block_on; 4 use radroots_event::SignedEvent; 5 use radroots_storage::{ 6 Error, 7 atomic::AtomicCommitDisposition, 8 authored::{ 9 AdmissionState, AuthoredArtifact, FailureClass, SigningState, WorkClaim, WorkFailure, 10 WorkPhase, 11 }, 12 authored_atomic::{ 13 ApplySignedArtifact, ApplyWorkFailure, AuthoredAtomicCommand, AuthoredAtomicReceipt, 14 AuthoredAtomicStorage, AuthoredWorkTarget, CancelAuthoredTarget, CancelAuthoredWork, 15 ClaimAuthoredTarget, ClaimAuthoredWork, RecordSignedArtifact, WorkFence, 16 }, 17 authored_delivery::AuthoredDeliveryState, 18 event::SourceGeneration, 19 memory::MemoryStorage, 20 }; 21 22 #[path = "authored_signing/fixture.rs"] 23 mod fixture; 24 use fixture::*; 25 26 fn prepared() -> (MemoryStorage, SignedEvent, WorkClaim, AuthoredAtomicReceipt) { 27 let storage = MemoryStorage::new(SourceGeneration::new([1; 32]).unwrap()); 28 let (preparation, event) = prepare(); 29 block_on(storage.execute_authored(preparation)).unwrap(); 30 let active = claim(NonZeroU64::new(1).unwrap(), 4, 11); 31 let receipt = block_on(storage.execute_authored(AuthoredAtomicCommand::Claim( 32 ClaimAuthoredWork::new( 33 ClaimAuthoredTarget::ArtifactSigning(ids().1), 34 active.clone(), 35 ), 36 ))) 37 .unwrap(); 38 (storage, event, active, receipt) 39 } 40 41 fn artifact(storage: &MemoryStorage) -> AuthoredArtifact { 42 block_on(storage.authored_artifact(ids().1)) 43 .unwrap() 44 .unwrap() 45 } 46 47 fn fence(active: &WorkClaim) -> WorkFence { 48 WorkFence::new(*active.token(), active.generation(), active.row_revision()).unwrap() 49 } 50 51 #[test] 52 fn expired_claim_retains_exact_facts_without_relaxing_active_fences() { 53 let (storage, event, active, original) = prepared(); 54 let stale = AuthoredAtomicCommand::ApplySigned( 55 ApplySignedArtifact::new(ids().1, fence(&active), event.clone(), 40).unwrap(), 56 ); 57 assert_eq!( 58 block_on(storage.execute_authored(stale.clone())), 59 Err(Error::DeliveryPlanClaimConflict) 60 ); 61 let fact = record(event.clone(), active.clone(), 40); 62 assert_ne!(fact.commit_id(), stale.commit_id()); 63 let receipt = block_on(storage.execute_authored(fact.clone())).unwrap(); 64 assert_eq!(receipt.disposition(), AtomicCommitDisposition::Committed); 65 let signed = artifact(&storage); 66 assert_eq!(signed.signing_state(), SigningState::Signed); 67 assert_eq!(signed.signed().unwrap().event().raw_json(), RAW); 68 assert!(signed.signing_claim().is_none()); 69 let delivery = block_on(storage.authored_delivery_plan(ids().2)) 70 .unwrap() 71 .unwrap(); 72 assert_eq!(delivery.request().unwrap().payload().event(), &event); 73 assert!(delivery.attempts().is_empty()); 74 assert_eq!( 75 block_on(storage.authored_receipt(original.commit_id())) 76 .unwrap() 77 .unwrap(), 78 original 79 ); 80 let later = record(event, active, 90); 81 assert_eq!(later.commit_id(), fact.commit_id()); 82 let replay = block_on(storage.execute_authored(later)).unwrap(); 83 assert_eq!(replay.disposition(), AtomicCommitDisposition::Replay); 84 assert_eq!(replay.committed_at_unix_ms(), 40); 85 assert_eq!(artifact(&storage), signed); 86 } 87 88 #[test] 89 fn cancellation_and_terminal_failure_survive_late_signed_bytes() { 90 for cancelled in [false, true] { 91 let (storage, event, active, _) = prepared(); 92 let stop = if cancelled { 93 AuthoredAtomicCommand::Cancel( 94 CancelAuthoredWork::new( 95 CancelAuthoredTarget::ArtifactSigning(ids().1), 96 artifact(&storage).revision(), 97 20, 98 ) 99 .unwrap(), 100 ) 101 } else { 102 AuthoredAtomicCommand::ApplyFailure( 103 ApplyWorkFailure::new( 104 AuthoredWorkTarget::Artifact(ids().1), 105 fence(&active), 106 WorkFailure::new( 107 "signing_stopped", 108 WorkPhase::Signing, 109 FailureClass::Terminal, 110 None, 111 None, 112 ) 113 .unwrap(), 114 None, 115 20, 116 ) 117 .unwrap(), 118 ) 119 }; 120 let stop_receipt = block_on(storage.execute_authored(stop)).unwrap(); 121 let stopped = artifact(&storage); 122 // An older observation still records the fact without moving row time back. 123 let fact_receipt = block_on(storage.execute_authored(record(event, active, 15))).unwrap(); 124 assert_eq!(fact_receipt.committed_at_unix_ms(), 20); 125 let retained = artifact(&storage); 126 assert_eq!(retained.signing_state(), stopped.signing_state()); 127 assert_eq!(retained.last_failure(), stopped.last_failure()); 128 assert_eq!(retained.updated_at_unix_ms(), 20); 129 assert_eq!(retained.admission_state(), AdmissionState::Pending); 130 assert_eq!(retained.signed().unwrap().event().raw_json(), RAW); 131 let delivery = block_on(storage.authored_delivery_plan(ids().2)) 132 .unwrap() 133 .unwrap(); 134 assert!(delivery.request().is_none()); 135 for target in [ 136 ClaimAuthoredTarget::ArtifactSigning(ids().1), 137 ClaimAuthoredTarget::ArtifactAdmission(ids().1), 138 ClaimAuthoredTarget::DeliveryPlan(ids().2), 139 ] { 140 let revision = if matches!(target, ClaimAuthoredTarget::DeliveryPlan(_)) { 141 delivery.revision() 142 } else { 143 retained.revision() 144 }; 145 assert!( 146 block_on(storage.execute_authored(AuthoredAtomicCommand::Claim( 147 ClaimAuthoredWork::new(target, claim(revision, 8, 50)), 148 ))) 149 .is_err() 150 ); 151 } 152 assert_eq!(artifact(&storage), retained); 153 assert_eq!( 154 block_on(storage.authored_receipt(stop_receipt.commit_id())) 155 .unwrap() 156 .unwrap(), 157 stop_receipt 158 ); 159 let operation = block_on(storage.authored_operation(ids().0)) 160 .unwrap() 161 .unwrap(); 162 let settlement = radroots_storage::authored::OperationSettlement::evaluate( 163 &operation, 164 std::slice::from_ref(&retained), 165 ) 166 .unwrap(); 167 assert_eq!(settlement.signed(), 1); 168 assert_eq!(settlement.cancelled(), u16::from(cancelled)); 169 assert_eq!(settlement.failed_terminal(), u16::from(!cancelled)); 170 #[cfg(feature = "serde")] 171 assert_eq!( 172 serde_json::from_str::<AuthoredArtifact>(&serde_json::to_string(&retained).unwrap()) 173 .unwrap(), 174 retained 175 ); 176 } 177 } 178 179 #[test] 180 fn superseded_attempts_cannot_overwrite_first_exact_bytes_or_restart_stopped_delivery() { 181 let (storage, event, first, _) = prepared(); 182 let second = claim(artifact(&storage).revision(), 5, 40); 183 block_on( 184 storage.execute_authored(AuthoredAtomicCommand::Claim(ClaimAuthoredWork::new( 185 ClaimAuthoredTarget::ArtifactSigning(ids().1), 186 second.clone(), 187 ))), 188 ) 189 .unwrap(); 190 let plan = block_on(storage.authored_delivery_plan(ids().2)) 191 .unwrap() 192 .unwrap(); 193 block_on( 194 storage.execute_authored(AuthoredAtomicCommand::Cancel( 195 CancelAuthoredWork::new( 196 CancelAuthoredTarget::DeliveryPlan(ids().2), 197 plan.revision(), 198 41, 199 ) 200 .unwrap(), 201 )), 202 ) 203 .unwrap(); 204 block_on(storage.execute_authored(record(event.clone(), first.clone(), 42))).unwrap(); 205 let retained = artifact(&storage); 206 assert!(retained.signing_claim().is_none()); 207 block_on(storage.execute_authored(record(event, second.clone(), 43))).unwrap(); 208 assert_eq!(artifact(&storage), retained); 209 let alternate = fixture::event(&format!(" {RAW} ")); 210 assert_eq!(alternate.id(), retained.signed().unwrap().event().id()); 211 assert_eq!( 212 block_on(storage.execute_authored(record(alternate, second, 44))), 213 Err(Error::AtomicCommitConflict) 214 ); 215 assert_eq!(artifact(&storage), retained); 216 let plan = block_on(storage.authored_delivery_plan(ids().2)) 217 .unwrap() 218 .unwrap(); 219 assert_eq!(plan.state(), AuthoredDeliveryState::Cancelled); 220 assert!(plan.request().is_none()); 221 assert!(plan.attempts().is_empty()); 222 } 223 224 #[test] 225 fn unrelated_plan_or_attempt_provenance_rolls_back_without_fact_receipts() { 226 let (storage, event, active, original) = prepared(); 227 let before = artifact(&storage); 228 let wrong_claims = [ 229 WorkClaim::new( 230 [9; 16], 231 active.owner(), 232 active.generation(), 233 11, 234 31, 235 active.row_revision(), 236 ) 237 .unwrap(), 238 WorkClaim::new( 239 *active.token(), 240 "different-worker", 241 active.generation(), 242 11, 243 31, 244 active.row_revision(), 245 ) 246 .unwrap(), 247 WorkClaim::new( 248 *active.token(), 249 active.owner(), 250 NonZeroU64::new(9).unwrap(), 251 11, 252 31, 253 active.row_revision(), 254 ) 255 .unwrap(), 256 WorkClaim::new( 257 *active.token(), 258 active.owner(), 259 active.generation(), 260 12, 261 31, 262 active.row_revision(), 263 ) 264 .unwrap(), 265 WorkClaim::new( 266 *active.token(), 267 active.owner(), 268 active.generation(), 269 11, 270 32, 271 active.row_revision(), 272 ) 273 .unwrap(), 274 WorkClaim::new( 275 *active.token(), 276 active.owner(), 277 active.generation(), 278 11, 279 31, 280 NonZeroU64::new(9).unwrap(), 281 ) 282 .unwrap(), 283 ]; 284 let expected = record(event.clone(), active.clone(), 50); 285 for wrong in wrong_claims { 286 let command = record(event.clone(), wrong, 50); 287 assert_ne!(command.commit_id(), expected.commit_id()); 288 assert!(block_on(storage.execute_authored(command.clone())).is_err()); 289 assert!( 290 block_on(storage.authored_receipt(command.commit_id())) 291 .unwrap() 292 .is_none() 293 ); 294 } 295 let wrong_operation = RecordSignedArtifact::new( 296 radroots_storage::journal::OperationInstanceId::new([9; 16]).unwrap(), 297 ids().1, 298 active.clone(), 299 event.clone(), 300 50, 301 ) 302 .unwrap(); 303 assert!( 304 block_on(storage.execute_authored(AuthoredAtomicCommand::RecordSigned(wrong_operation))) 305 .is_err() 306 ); 307 let mismatch = record(fixture::event(OTHER_RAW), active.clone(), 50); 308 assert!(block_on(storage.execute_authored(mismatch.clone())).is_err()); 309 assert!( 310 block_on(storage.authored_receipt(mismatch.commit_id())) 311 .unwrap() 312 .is_none() 313 ); 314 assert_eq!(artifact(&storage), before); 315 let AuthoredAtomicCommand::RecordSigned(value) = expected else { 316 unreachable!() 317 }; 318 assert_eq!(value.operation_id(), ids().0); 319 assert_eq!(value.artifact_id(), ids().1); 320 assert_eq!(value.claim(), &active); 321 assert_eq!(value.event(), &event); 322 assert_eq!(value.observed_at_unix_ms(), 50); 323 assert_eq!(value.clone(), value); 324 let (preparation, _) = prepare(); 325 let wrong_outcome = block_on(storage.execute_authored(preparation)).unwrap(); 326 assert!(value.apply_to(&mut before.clone(), &wrong_outcome).is_err()); 327 assert!(value.apply_to(&mut before.clone(), &original).is_ok()); 328 #[cfg(feature = "serde")] 329 { 330 let mut regressed = serde_json::to_value(&before).unwrap(); 331 regressed["signing_claim"] = serde_json::Value::Null; 332 regressed["updated_at_unix_ms"] = serde_json::Value::from(10); 333 let mut regressed = serde_json::from_value::<AuthoredArtifact>(regressed).unwrap(); 334 assert!(value.apply_to(&mut regressed, &original).is_err()); 335 } 336 } 337 338 #[test] 339 fn invalid_signature_and_pre_attempt_observation_never_become_facts() { 340 let (_, event, active, _) = prepared(); 341 for at in [0, 10] { 342 assert!( 343 RecordSignedArtifact::new(ids().0, ids().1, active.clone(), event.clone(), at).is_err() 344 ); 345 } 346 let mut value: serde_json::Value = serde_json::from_str(RAW).unwrap(); 347 value["sig"] = serde_json::Value::String("ff".repeat(64)); 348 let hostile = fixture::event(&serde_json::to_string(&value).unwrap()); 349 assert_eq!( 350 RecordSignedArtifact::new(ids().0, ids().1, active, hostile, 50), 351 Err(Error::InvalidAuthoredArtifact) 352 ); 353 } 354 355 #[test] 356 fn indeterminate_signing_is_resolved_by_original_verified_evidence() { 357 let (storage, event, active, _) = prepared(); 358 block_on( 359 storage.execute_authored(AuthoredAtomicCommand::ApplyFailure( 360 ApplyWorkFailure::new( 361 AuthoredWorkTarget::Artifact(ids().1), 362 fence(&active), 363 WorkFailure::new( 364 "signing_unknown", 365 WorkPhase::Signing, 366 FailureClass::Indeterminate, 367 None, 368 None, 369 ) 370 .unwrap(), 371 None, 372 20, 373 ) 374 .unwrap(), 375 )), 376 ) 377 .unwrap(); 378 assert_eq!( 379 artifact(&storage).signing_state(), 380 SigningState::Indeterminate 381 ); 382 block_on(storage.execute_authored(record(event, active, 40))).unwrap(); 383 let retained = artifact(&storage); 384 assert_eq!(retained.signing_state(), SigningState::Signed); 385 assert!(retained.last_failure().is_none()); 386 assert_eq!(retained.signed().unwrap().event().raw_json(), RAW); 387 } 388 389 #[test] 390 fn signed_fact_receipts_require_exact_outcome_and_monotonic_commit_time() { 391 let (storage, event, active, _) = prepared(); 392 let command = record(event, active, 40); 393 let unsigned = artifact(&storage); 394 assert!( 395 AuthoredAtomicReceipt::new( 396 &command, 397 AtomicCommitDisposition::Committed, 398 40, 399 radroots_storage::authored_atomic::AuthoredAtomicOutcome::Artifact(unsigned), 400 ) 401 .is_err() 402 ); 403 let receipt = block_on(storage.execute_authored(command.clone())).unwrap(); 404 assert!(receipt.matches_command(&command)); 405 assert!( 406 AuthoredAtomicReceipt::new( 407 &command, 408 AtomicCommitDisposition::Committed, 409 39, 410 receipt.outcome().clone(), 411 ) 412 .is_err() 413 ); 414 let wrong = radroots_storage::authored_atomic::AuthoredAtomicOutcome::DeliveryPlan( 415 block_on(storage.authored_delivery_plan(ids().2)) 416 .unwrap() 417 .unwrap(), 418 ); 419 assert!( 420 AuthoredAtomicReceipt::new( 421 &command, 422 AtomicCommitDisposition::Committed, 423 40, 424 wrong.clone() 425 ) 426 .is_err() 427 ); 428 let malformed = AuthoredAtomicReceipt::from_durable_parts( 429 command.commit_id(), 430 command.digest(), 431 AtomicCommitDisposition::Committed, 432 40, 433 wrong, 434 ) 435 .unwrap(); 436 assert!(!malformed.matches_command(&command)); 437 } 438 439 #[cfg(feature = "serde")] 440 #[test] 441 fn stopped_signed_snapshots_cannot_erase_the_required_terminal_failure() { 442 let (storage, event, active, _) = prepared(); 443 block_on( 444 storage.execute_authored(AuthoredAtomicCommand::ApplyFailure( 445 ApplyWorkFailure::new( 446 AuthoredWorkTarget::Artifact(ids().1), 447 fence(&active), 448 WorkFailure::new( 449 "signing_stopped", 450 WorkPhase::Signing, 451 FailureClass::Terminal, 452 None, 453 None, 454 ) 455 .unwrap(), 456 None, 457 20, 458 ) 459 .unwrap(), 460 )), 461 ) 462 .unwrap(); 463 block_on(storage.execute_authored(record(event, active, 40))).unwrap(); 464 let retained = artifact(&storage); 465 let mut forged = serde_json::to_value(&retained).unwrap(); 466 forged["last_failure"] = serde_json::Value::Null; 467 assert!(serde_json::from_value::<AuthoredArtifact>(forged).is_err()); 468 }