authored.rs (41229B)
1 use core::num::{NonZeroU32, NonZeroU64}; 2 use radroots_event::{GenericEventDraft, SignedEvent, wire::v1::Nip01EventWire}; 3 use radroots_event_codec::authoring::AuthoredEventPlan; 4 use radroots_storage::{ 5 Error, 6 authored::{ 7 AUTHORED_OPERATION_ARTIFACTS_MAX, AdmissionState, ArtifactOrigin, AuthoredArtifact, 8 AuthoredArtifactId, AuthoredOperation, FailureClass, OperationSettlement, RetrySchedule, 9 SigningState, WORK_CLAIM_OWNER_MAX_BYTES, WORK_FAILURE_CODE_MAX_BYTES, 10 WORK_FAILURE_DIAGNOSTIC_MAX_BYTES, WorkClaim, WorkFailure, WorkPhase, 11 }, 12 authored_delivery::{AuthoredDeliveryPlan, AuthoredDeliveryPlanId}, 13 journal::OperationInstanceId, 14 }; 15 use radroots_transport::{ 16 DeliveryReceipt, DeliveryRequest, SinkFailure, Target, TargetSet, 17 outcome::{DeliveryOutcome, Retryability}, 18 policy::{SatisfactionClass, SatisfactionPolicy, TargetPolicy}, 19 sink::{DeliveryPayload, DeliveryTargetReceipt}, 20 }; 21 22 const AUTHOR: &str = "585591529da0bab31b3b1b1f986611cf5f435dca84f978c89ee8a40cca7103df"; 23 const CREATED_AT: u64 = 1_800_000_100; 24 25 fn operation_id() -> OperationInstanceId { 26 OperationInstanceId::new([1; 16]).expect("operation ID") 27 } 28 29 fn artifact_id(value: u8) -> AuthoredArtifactId { 30 AuthoredArtifactId::new([value; 16]).expect("artifact ID") 31 } 32 33 fn plan() -> AuthoredEventPlan { 34 AuthoredEventPlan::from_generic( 35 GenericEventDraft::new( 36 "radroots.social.geochat.v1", 37 20_000, 38 CREATED_AT, 39 Vec::new(), 40 "durable authored plan", 41 AUTHOR, 42 ) 43 .expect("draft"), 44 ) 45 .expect("plan") 46 } 47 48 fn signed(plan: &AuthoredEventPlan) -> SignedEvent { 49 let wire = Nip01EventWire { 50 id: plan.expected_event_id().to_hex(), 51 pubkey: plan.author().to_hex(), 52 created_at: plan.created_at(), 53 kind: plan.body().kind(), 54 tags: plan.body().tags().to_vec(), 55 content: plan.body().content().to_owned(), 56 sig: "dd".repeat(64), 57 extra: Default::default(), 58 }; 59 let raw = serde_json::to_string(&wire).expect("raw event"); 60 SignedEvent::from_wire_verified_id(wire, raw).expect("signed event") 61 } 62 63 fn retry_failure(phase: WorkPhase, at: u64) -> WorkFailure { 64 WorkFailure::new( 65 "temporary_failure", 66 phase, 67 FailureClass::Retryable, 68 Some(at), 69 Some("temporary failure".to_owned()), 70 ) 71 .expect("failure") 72 } 73 74 #[test] 75 fn claim_failure_and_retry_models_are_bounded_and_fenced() { 76 let revision = NonZeroU64::MIN; 77 let claim = 78 WorkClaim::new([2; 16], "worker-1", NonZeroU64::MIN, 10, 20, revision).expect("claim"); 79 assert!(claim.matches_fence(&[2; 16], NonZeroU64::MIN, revision, 10)); 80 assert!(!claim.matches_fence(&[2; 16], NonZeroU64::MIN, revision, 20)); 81 assert_eq!( 82 WorkClaim::new([0; 16], "worker", NonZeroU64::MIN, 10, 20, revision), 83 Err(Error::InvalidWorkClaim) 84 ); 85 86 let failure = retry_failure(WorkPhase::Signing, 30); 87 let retry = RetrySchedule::new(NonZeroU32::MIN, 30, failure.clone()).expect("retry"); 88 assert_eq!(retry.attempt().get(), 1); 89 assert_eq!(retry.failure(), &failure); 90 assert_eq!( 91 WorkFailure::new( 92 "INVALID", 93 WorkPhase::Signing, 94 FailureClass::Terminal, 95 None, 96 None, 97 ), 98 Err(Error::InvalidWorkFailure) 99 ); 100 assert_eq!( 101 RetrySchedule::new(NonZeroU32::MIN, 31, retry_failure(WorkPhase::Signing, 30),), 102 Err(Error::InvalidRetrySchedule) 103 ); 104 } 105 106 #[test] 107 fn planned_artifacts_enforce_exact_signing_and_admission_transitions() { 108 let plan = plan(); 109 let mut artifact = AuthoredArtifact::planned(artifact_id(2), operation_id(), 0, &plan, 10) 110 .expect("planned artifact"); 111 let claim = WorkClaim::new( 112 [3; 16], 113 "signer", 114 NonZeroU64::MIN, 115 11, 116 20, 117 artifact.revision(), 118 ) 119 .expect("claim"); 120 artifact 121 .set_signing_claim(claim, 11) 122 .expect("claim signing"); 123 artifact.record_signed(signed(&plan), 12).expect("sign"); 124 assert_eq!(artifact.signing_state(), SigningState::Signed); 125 assert!(artifact.signed().is_some()); 126 assert!(artifact.signing_claim().is_none()); 127 artifact 128 .record_admission(AdmissionState::Duplicate, None, None, 13) 129 .expect("admit duplicate"); 130 assert!(artifact.admission_state().is_admitted()); 131 132 let other = AuthoredEventPlan::from_generic( 133 GenericEventDraft::new( 134 "radroots.social.geochat.v1", 135 20_000, 136 CREATED_AT, 137 Vec::new(), 138 "different", 139 AUTHOR, 140 ) 141 .expect("other draft"), 142 ) 143 .expect("other plan"); 144 let mut mismatched = AuthoredArtifact::planned(artifact_id(3), operation_id(), 1, &plan, 10) 145 .expect("planned artifact"); 146 assert_eq!( 147 mismatched.record_signed(signed(&other), 11), 148 Err(Error::InvalidAuthoredArtifact) 149 ); 150 assert_eq!(mismatched.signing_state(), SigningState::Planned); 151 assert!(mismatched.signed().is_none()); 152 } 153 154 #[test] 155 fn imported_artifacts_are_non_resignable_and_settlement_preserves_order() { 156 let plan = plan(); 157 let imported = 158 AuthoredArtifact::imported_signed(artifact_id(2), operation_id(), 0, signed(&plan), 10) 159 .expect("imported"); 160 assert_eq!(imported.origin(), ArtifactOrigin::ImportedSigned); 161 assert!(!imported.origin().is_resignable()); 162 163 let mut planned = 164 AuthoredArtifact::planned(artifact_id(3), operation_id(), 1, &plan, 10).expect("planned"); 165 planned.cancel_signing(11).expect("cancel"); 166 let operation = 167 AuthoredOperation::new(operation_id(), vec![artifact_id(2), artifact_id(3)], 10) 168 .expect("operation"); 169 let settlement = 170 OperationSettlement::evaluate(&operation, &[imported, planned]).expect("settlement"); 171 assert_eq!(settlement.artifacts(), 2); 172 assert_eq!(settlement.signed(), 1); 173 assert_eq!(settlement.pending(), 1); 174 assert_eq!(settlement.cancelled(), 1); 175 assert!(!settlement.is_settled()); 176 177 assert_eq!( 178 AuthoredOperation::new(operation_id(), vec![artifact_id(2), artifact_id(2)], 10,), 179 Err(Error::InvalidAuthoredOperation) 180 ); 181 } 182 183 #[test] 184 fn durable_models_round_trip_and_reject_forged_state() { 185 let plan = plan(); 186 let mut artifact = 187 AuthoredArtifact::planned(artifact_id(2), operation_id(), 0, &plan, 10).expect("planned"); 188 let failure = retry_failure(WorkPhase::Signing, 20); 189 artifact 190 .record_signing_failure( 191 failure.clone(), 192 Some(RetrySchedule::new(NonZeroU32::MIN, 20, failure).expect("retry")), 193 11, 194 ) 195 .expect("retryable signing"); 196 let json = serde_json::to_string(&artifact).expect("artifact json"); 197 assert_eq!( 198 serde_json::from_str::<AuthoredArtifact>(&json).expect("artifact round trip"), 199 artifact 200 ); 201 202 let mut forged: serde_json::Value = serde_json::from_str(&json).expect("artifact value"); 203 forged["signing_state"] = serde_json::json!("signed"); 204 assert!(serde_json::from_value::<AuthoredArtifact>(forged).is_err()); 205 206 let exact = artifact.plan().expect("durable plan"); 207 let mut corrupt = exact.wire_json().to_vec(); 208 corrupt[0] ^= 1; 209 assert_eq!( 210 radroots_storage::authored::DurableAuthoredPlan::reconstruct(corrupt), 211 Err(Error::InvalidAuthoredArtifact) 212 ); 213 214 let imported = 215 AuthoredArtifact::imported_signed(artifact_id(4), operation_id(), 0, signed(&plan), 10) 216 .expect("imported"); 217 let mut forged_digest = serde_json::to_value(imported).expect("imported value"); 218 forged_digest["signed"]["raw_json_sha256"] = 219 serde_json::to_value([0_u8; 32]).expect("digest value"); 220 assert!(serde_json::from_value::<AuthoredArtifact>(forged_digest).is_err()); 221 222 let operation = 223 AuthoredOperation::new(operation_id(), vec![artifact_id(2), artifact_id(3)], 10) 224 .expect("operation"); 225 let operation_json = serde_json::to_string(&operation).expect("operation json"); 226 assert_eq!( 227 serde_json::from_str::<AuthoredOperation>(&operation_json).expect("operation round trip"), 228 operation 229 ); 230 let mut forged_operation: serde_json::Value = 231 serde_json::from_str(&operation_json).expect("operation value"); 232 forged_operation["artifact_ids"][1] = forged_operation["artifact_ids"][0].clone(); 233 assert!(serde_json::from_value::<AuthoredOperation>(forged_operation).is_err()); 234 } 235 236 #[test] 237 fn claim_validation_and_fence_dimensions_are_independent() { 238 let revision = NonZeroU64::new(7).unwrap(); 239 let generation = NonZeroU64::new(3).unwrap(); 240 let claim = WorkClaim::new([9; 16], "worker", generation, 10, 20, revision).unwrap(); 241 assert_eq!(claim.token(), &[9; 16]); 242 assert_eq!(claim.owner(), "worker"); 243 assert_eq!(claim.generation(), generation); 244 assert_eq!(claim.acquired_at_unix_ms(), 10); 245 assert_eq!(claim.expires_at_unix_ms(), 20); 246 assert_eq!(claim.row_revision(), revision); 247 assert!(claim.matches_fence(&[9; 16], generation, revision, 10)); 248 assert!(claim.matches_fence(&[9; 16], generation, revision, 19)); 249 assert!(!claim.matches_fence(&[8; 16], generation, revision, 10)); 250 assert!(!claim.matches_fence(&[9; 16], NonZeroU64::MIN, revision, 10)); 251 assert!(!claim.matches_fence(&[9; 16], generation, NonZeroU64::MIN, 10)); 252 assert!(!claim.matches_fence(&[9; 16], generation, revision, 9)); 253 assert!(!claim.matches_fence(&[9; 16], generation, revision, 20)); 254 255 for (token, owner, acquired, expires) in [ 256 ([0; 16], "worker".to_owned(), 10, 20), 257 ([1; 16], String::new(), 10, 20), 258 ([1; 16], " worker".to_owned(), 10, 20), 259 ([1; 16], "worker ".to_owned(), 10, 20), 260 ([1; 16], "worker\nname".to_owned(), 10, 20), 261 ([1; 16], "x".repeat(WORK_CLAIM_OWNER_MAX_BYTES + 1), 10, 20), 262 ([1; 16], "worker".to_owned(), 0, 20), 263 ([1; 16], "worker".to_owned(), 10, 10), 264 ([1; 16], "worker".to_owned(), 10, 9), 265 ] { 266 assert_eq!( 267 WorkClaim::new(token, owner, generation, acquired, expires, revision), 268 Err(Error::InvalidWorkClaim) 269 ); 270 } 271 272 let json = serde_json::to_string(&claim).unwrap(); 273 assert_eq!(serde_json::from_str::<WorkClaim>(&json).unwrap(), claim); 274 let mut invalid: serde_json::Value = serde_json::from_str(&json).unwrap(); 275 invalid["owner"] = serde_json::json!(""); 276 assert!(serde_json::from_value::<WorkClaim>(invalid).is_err()); 277 } 278 279 #[test] 280 fn failure_and_retry_models_cover_every_class_and_boundary() { 281 for phase in [ 282 WorkPhase::Signing, 283 WorkPhase::Admission, 284 WorkPhase::Delivery, 285 ] { 286 for class in [ 287 FailureClass::Retryable, 288 FailureClass::Terminal, 289 FailureClass::Indeterminate, 290 ] { 291 let retry_after = (class == FailureClass::Retryable).then_some(20); 292 let failure = WorkFailure::new( 293 "failure.code-1", 294 phase, 295 class, 296 retry_after, 297 Some("safe detail".to_owned()), 298 ) 299 .unwrap(); 300 assert_eq!(failure.code(), "failure.code-1"); 301 assert_eq!(failure.phase(), phase); 302 assert_eq!(failure.class(), class); 303 assert_eq!(failure.retry_after_unix_ms(), retry_after); 304 assert_eq!(failure.diagnostic(), Some("safe detail")); 305 failure.validate().unwrap(); 306 assert_eq!( 307 serde_json::from_str::<WorkFailure>(&serde_json::to_string(&failure).unwrap()) 308 .unwrap(), 309 failure 310 ); 311 } 312 } 313 314 for (code, class, retry_after, diagnostic) in [ 315 ("", FailureClass::Terminal, None, None), 316 ("INVALID", FailureClass::Terminal, None, None), 317 ("bad/code", FailureClass::Terminal, None, None), 318 ( 319 "x".repeat(WORK_FAILURE_CODE_MAX_BYTES + 1).leak(), 320 FailureClass::Terminal, 321 None, 322 None, 323 ), 324 ("failure", FailureClass::Retryable, Some(0), None), 325 ("failure", FailureClass::Terminal, Some(20), None), 326 ("failure", FailureClass::Indeterminate, Some(20), None), 327 ("failure", FailureClass::Terminal, None, Some("".to_owned())), 328 ( 329 "failure", 330 FailureClass::Terminal, 331 None, 332 Some(" diagnostic".to_owned()), 333 ), 334 ( 335 "failure", 336 FailureClass::Terminal, 337 None, 338 Some("x".repeat(WORK_FAILURE_DIAGNOSTIC_MAX_BYTES + 1)), 339 ), 340 ] { 341 assert_eq!( 342 WorkFailure::new(code, WorkPhase::Signing, class, retry_after, diagnostic), 343 Err(Error::InvalidWorkFailure) 344 ); 345 } 346 347 let initial_failure = retry_failure(WorkPhase::Delivery, 30); 348 let retry = RetrySchedule::new(NonZeroU32::MIN, 30, initial_failure.clone()).unwrap(); 349 assert_eq!(retry.attempt(), NonZeroU32::MIN); 350 assert_eq!(retry.not_before_unix_ms(), 30); 351 assert_eq!(retry.failure(), &initial_failure); 352 let next_failure = retry_failure(WorkPhase::Delivery, 40); 353 assert_eq!( 354 retry 355 .next_attempt(40, next_failure) 356 .unwrap() 357 .attempt() 358 .get(), 359 2 360 ); 361 assert_eq!( 362 RetrySchedule::new(NonZeroU32::MIN, 0, initial_failure.clone()), 363 Err(Error::InvalidRetrySchedule) 364 ); 365 assert_eq!( 366 RetrySchedule::new( 367 NonZeroU32::MIN, 368 30, 369 WorkFailure::new( 370 "terminal", 371 WorkPhase::Delivery, 372 FailureClass::Terminal, 373 None, 374 None, 375 ) 376 .unwrap(), 377 ), 378 Err(Error::InvalidRetrySchedule) 379 ); 380 assert_eq!( 381 RetrySchedule::new(NonZeroU32::MIN, 31, initial_failure), 382 Err(Error::InvalidRetrySchedule) 383 ); 384 385 let mut maximum = serde_json::to_value(&retry).unwrap(); 386 maximum["attempt"] = serde_json::json!(u32::MAX); 387 let maximum: RetrySchedule = serde_json::from_value(maximum).unwrap(); 388 assert_eq!( 389 maximum.next_attempt(40, retry_failure(WorkPhase::Delivery, 40)), 390 Err(Error::InvalidRetrySchedule) 391 ); 392 } 393 394 #[test] 395 fn operation_construction_reconstruction_and_accessors_are_bounded() { 396 assert_eq!( 397 AuthoredArtifactId::new([0; 16]), 398 Err(Error::InvalidAuthoredArtifact) 399 ); 400 let id = artifact_id(2); 401 assert_eq!(id.as_bytes(), &[2; 16]); 402 assert_eq!(AuthoredArtifactId::try_from([2; 16]).unwrap(), id); 403 assert_eq!(<[u8; 16]>::from(id), [2; 16]); 404 405 let operation = AuthoredOperation::new(operation_id(), vec![id], 10).unwrap(); 406 assert_eq!(operation.operation_id(), operation_id()); 407 assert_eq!(operation.artifact_ids(), &[id]); 408 assert_eq!(operation.created_at_unix_ms(), 10); 409 assert_eq!(operation.updated_at_unix_ms(), 10); 410 assert_eq!(operation.revision(), NonZeroU64::MIN); 411 assert_eq!( 412 AuthoredOperation::reconstruct( 413 operation_id(), 414 vec![id], 415 10, 416 11, 417 NonZeroU64::new(2).unwrap(), 418 ) 419 .unwrap() 420 .updated_at_unix_ms(), 421 11 422 ); 423 424 for (ids, created, updated) in [ 425 (Vec::new(), 10, 10), 426 (vec![id, id], 10, 10), 427 ( 428 (0..=AUTHORED_OPERATION_ARTIFACTS_MAX) 429 .map(|index| artifact_id((index % 254 + 1) as u8)) 430 .collect(), 431 10, 432 10, 433 ), 434 (vec![id], 0, 10), 435 (vec![id], 10, 9), 436 ] { 437 assert_eq!( 438 AuthoredOperation::reconstruct(operation_id(), ids, created, updated, NonZeroU64::MIN,), 439 Err(Error::InvalidAuthoredOperation) 440 ); 441 } 442 } 443 444 fn claim_for(artifact: &AuthoredArtifact, token: u8, generation: u64, acquired: u64) -> WorkClaim { 445 WorkClaim::new( 446 [token; 16], 447 format!("worker-{token}"), 448 NonZeroU64::new(generation).unwrap(), 449 acquired, 450 acquired + 10, 451 artifact.revision(), 452 ) 453 .unwrap() 454 } 455 456 fn failure(phase: WorkPhase, class: FailureClass, at: Option<u64>) -> WorkFailure { 457 WorkFailure::new("operation_failure", phase, class, at, None).unwrap() 458 } 459 460 #[test] 461 fn signing_claim_failure_terminal_indeterminate_and_cancel_paths_are_fenced() { 462 let authored_plan = plan(); 463 let mut retryable = 464 AuthoredArtifact::planned(artifact_id(10), operation_id(), 0, &authored_plan, 10).unwrap(); 465 let first_claim = claim_for(&retryable, 1, 1, 11); 466 retryable 467 .set_signing_claim(first_claim.clone(), 11) 468 .unwrap(); 469 assert_eq!(retryable.signing_claim(), Some(&first_claim)); 470 assert_eq!( 471 retryable.set_signing_claim(claim_for(&retryable, 2, 2, 12), 12), 472 Err(Error::InvalidAuthoredTransition) 473 ); 474 let replacement = claim_for(&retryable, 3, 2, 21); 475 retryable 476 .set_signing_claim(replacement.clone(), 21) 477 .unwrap(); 478 assert_eq!(retryable.signing_claim(), Some(&replacement)); 479 480 let retry_failure = failure(WorkPhase::Signing, FailureClass::Retryable, Some(40)); 481 let retry = RetrySchedule::new(NonZeroU32::MIN, 40, retry_failure.clone()).unwrap(); 482 retryable 483 .record_signing_failure(retry_failure.clone(), Some(retry.clone()), 22) 484 .unwrap(); 485 assert_eq!(retryable.signing_state(), SigningState::Retryable); 486 assert_eq!(retryable.signing_retry(), Some(&retry)); 487 assert_eq!(retryable.last_failure(), Some(&retry_failure)); 488 assert_eq!( 489 retryable.set_signing_claim(claim_for(&retryable, 4, 3, 39), 39), 490 Err(Error::InvalidAuthoredTransition) 491 ); 492 let retry_claim = claim_for(&retryable, 4, 3, 40); 493 retryable.set_signing_claim(retry_claim, 40).unwrap(); 494 retryable.record_signed(signed(&authored_plan), 41).unwrap(); 495 assert_eq!(retryable.signing_state(), SigningState::Signed); 496 assert_eq!(retryable.signing_retry(), None); 497 assert_eq!(retryable.last_failure(), None); 498 499 let mut terminal = 500 AuthoredArtifact::planned(artifact_id(11), operation_id(), 0, &authored_plan, 10).unwrap(); 501 let terminal_failure = failure(WorkPhase::Signing, FailureClass::Terminal, None); 502 terminal 503 .record_signing_failure(terminal_failure.clone(), None, 11) 504 .unwrap(); 505 assert_eq!(terminal.signing_state(), SigningState::FailedTerminal); 506 assert_eq!(terminal.last_failure(), Some(&terminal_failure)); 507 assert_eq!( 508 terminal.cancel_signing(12), 509 Err(Error::InvalidAuthoredTransition) 510 ); 511 512 let mut indeterminate = 513 AuthoredArtifact::planned(artifact_id(12), operation_id(), 0, &authored_plan, 10).unwrap(); 514 let indeterminate_failure = failure(WorkPhase::Signing, FailureClass::Indeterminate, None); 515 indeterminate 516 .record_signing_failure(indeterminate_failure.clone(), None, 11) 517 .unwrap(); 518 assert_eq!(indeterminate.signing_state(), SigningState::Indeterminate); 519 520 let mut cancelled = 521 AuthoredArtifact::planned(artifact_id(13), operation_id(), 0, &authored_plan, 10).unwrap(); 522 cancelled.cancel_signing(11).unwrap(); 523 assert_eq!(cancelled.signing_state(), SigningState::Cancelled); 524 assert_eq!(cancelled.updated_at_unix_ms(), 11); 525 526 let mut imported = AuthoredArtifact::imported_signed( 527 artifact_id(14), 528 operation_id(), 529 0, 530 signed(&authored_plan), 531 10, 532 ) 533 .unwrap(); 534 assert_eq!( 535 imported.record_signed(signed(&authored_plan), 11), 536 Err(Error::InvalidAuthoredTransition) 537 ); 538 assert_eq!( 539 imported.cancel_signing(11), 540 Err(Error::InvalidAuthoredTransition) 541 ); 542 543 let mut invalid = 544 AuthoredArtifact::planned(artifact_id(15), operation_id(), 0, &authored_plan, 10).unwrap(); 545 assert_eq!( 546 invalid.record_signing_failure( 547 failure(WorkPhase::Admission, FailureClass::Terminal, None), 548 None, 549 11, 550 ), 551 Err(Error::InvalidAuthoredTransition) 552 ); 553 assert_eq!( 554 invalid.record_signing_failure( 555 failure(WorkPhase::Signing, FailureClass::Retryable, Some(20)), 556 None, 557 11, 558 ), 559 Err(Error::InvalidAuthoredTransition) 560 ); 561 assert_eq!( 562 invalid.record_signing_failure( 563 failure(WorkPhase::Signing, FailureClass::Terminal, None), 564 Some( 565 RetrySchedule::new( 566 NonZeroU32::MIN, 567 20, 568 failure(WorkPhase::Signing, FailureClass::Retryable, Some(20)), 569 ) 570 .unwrap(), 571 ), 572 11, 573 ), 574 Err(Error::InvalidAuthoredTransition) 575 ); 576 let before = invalid.clone(); 577 assert_eq!( 578 invalid.cancel_signing(9), 579 Err(Error::InvalidAuthoredTransition) 580 ); 581 assert_eq!(invalid, before); 582 } 583 584 fn signed_artifact(value: u8) -> AuthoredArtifact { 585 let authored_plan = plan(); 586 let mut artifact = 587 AuthoredArtifact::planned(artifact_id(value), operation_id(), 0, &authored_plan, 10) 588 .unwrap(); 589 artifact.record_signed(signed(&authored_plan), 11).unwrap(); 590 artifact 591 } 592 593 #[test] 594 fn admission_claim_and_result_paths_enforce_failure_coherence() { 595 let mut inserted = signed_artifact(20); 596 let active = claim_for(&inserted, 1, 1, 12); 597 inserted.set_admission_claim(active.clone(), 12).unwrap(); 598 assert_eq!(inserted.admission_claim(), Some(&active)); 599 inserted 600 .record_admission(AdmissionState::Inserted, None, None, 13) 601 .unwrap(); 602 assert_eq!(inserted.admission_state(), AdmissionState::Inserted); 603 assert!(inserted.admission_state().is_admitted()); 604 assert_eq!(inserted.admission_claim(), None); 605 606 let mut duplicate = signed_artifact(21); 607 duplicate 608 .record_admission(AdmissionState::Duplicate, None, None, 12) 609 .unwrap(); 610 assert!(duplicate.admission_state().is_admitted()); 611 612 let mut retryable = signed_artifact(22); 613 let retry_failure = failure(WorkPhase::Admission, FailureClass::Retryable, Some(30)); 614 let retry = RetrySchedule::new(NonZeroU32::MIN, 30, retry_failure.clone()).unwrap(); 615 retryable 616 .record_admission( 617 AdmissionState::Retryable, 618 Some(retry_failure.clone()), 619 Some(retry.clone()), 620 12, 621 ) 622 .unwrap(); 623 assert_eq!(retryable.admission_retry(), Some(&retry)); 624 assert_eq!(retryable.last_failure(), Some(&retry_failure)); 625 assert_eq!( 626 retryable.set_admission_claim(claim_for(&retryable, 2, 2, 29), 29), 627 Err(Error::InvalidAuthoredTransition) 628 ); 629 let active = claim_for(&retryable, 2, 2, 30); 630 retryable.set_admission_claim(active, 30).unwrap(); 631 632 for (value, state) in [ 633 (23, AdmissionState::Rejected), 634 (24, AdmissionState::Cancelled), 635 ] { 636 let mut artifact = signed_artifact(value); 637 let terminal = failure(WorkPhase::Admission, FailureClass::Terminal, None); 638 artifact 639 .record_admission(state, Some(terminal.clone()), None, 12) 640 .unwrap(); 641 assert_eq!(artifact.admission_state(), state); 642 assert_eq!(artifact.last_failure(), Some(&terminal)); 643 } 644 645 let invalid_cases = [ 646 (AdmissionState::Pending, None, None), 647 ( 648 AdmissionState::Inserted, 649 Some(failure(WorkPhase::Admission, FailureClass::Terminal, None)), 650 None, 651 ), 652 (AdmissionState::Rejected, None, None), 653 ( 654 AdmissionState::Retryable, 655 Some(failure( 656 WorkPhase::Admission, 657 FailureClass::Retryable, 658 Some(30), 659 )), 660 None, 661 ), 662 ( 663 AdmissionState::Rejected, 664 Some(failure(WorkPhase::Signing, FailureClass::Terminal, None)), 665 None, 666 ), 667 ]; 668 for (index, (state, failure, retry)) in invalid_cases.into_iter().enumerate() { 669 let mut artifact = signed_artifact(30 + index as u8); 670 assert_eq!( 671 artifact.record_admission(state, failure, retry, 12), 672 Err(Error::InvalidAuthoredTransition) 673 ); 674 } 675 676 let mut unsigned = 677 AuthoredArtifact::planned(artifact_id(40), operation_id(), 0, &plan(), 10).unwrap(); 678 assert_eq!( 679 unsigned.set_admission_claim(claim_for(&unsigned, 1, 1, 11), 11), 680 Err(Error::InvalidAuthoredTransition) 681 ); 682 } 683 684 fn assert_artifact_json_rejected( 685 mut value: serde_json::Value, 686 key: &str, 687 replacement: serde_json::Value, 688 ) { 689 value[key] = replacement; 690 assert!( 691 serde_json::from_value::<AuthoredArtifact>(value).is_err(), 692 "forged artifact field {key} was accepted" 693 ); 694 } 695 696 #[test] 697 fn authored_artifact_reconstruction_rejects_each_incoherent_durable_state() { 698 let authored_plan = plan(); 699 let planned = 700 AuthoredArtifact::planned(artifact_id(50), operation_id(), 0, &authored_plan, 10).unwrap(); 701 let planned_value = serde_json::to_value(&planned).unwrap(); 702 assert_artifact_json_rejected( 703 planned_value.clone(), 704 "created_at_unix_ms", 705 serde_json::json!(0), 706 ); 707 assert_artifact_json_rejected( 708 planned_value.clone(), 709 "updated_at_unix_ms", 710 serde_json::json!(9), 711 ); 712 assert_artifact_json_rejected(planned_value.clone(), "plan", serde_json::Value::Null); 713 assert_artifact_json_rejected( 714 planned_value.clone(), 715 "signing_state", 716 serde_json::json!("signed"), 717 ); 718 assert_artifact_json_rejected( 719 planned_value.clone(), 720 "admission_state", 721 serde_json::json!("inserted"), 722 ); 723 assert_artifact_json_rejected( 724 planned_value.clone(), 725 "signing_state", 726 serde_json::json!("retryable"), 727 ); 728 let mut admission_without_signed = planned_value.clone(); 729 admission_without_signed["admission_state"] = serde_json::json!("retryable"); 730 assert!(serde_json::from_value::<AuthoredArtifact>(admission_without_signed).is_err()); 731 732 let imported = AuthoredArtifact::imported_signed( 733 artifact_id(51), 734 operation_id(), 735 0, 736 signed(&authored_plan), 737 10, 738 ) 739 .unwrap(); 740 let imported_value = serde_json::to_value(&imported).unwrap(); 741 let mut imported_with_plan = imported_value.clone(); 742 imported_with_plan["plan"] = planned_value["plan"].clone(); 743 assert!(serde_json::from_value::<AuthoredArtifact>(imported_with_plan).is_err()); 744 assert_artifact_json_rejected( 745 imported_value.clone(), 746 "signing_state", 747 serde_json::json!("planned"), 748 ); 749 assert_artifact_json_rejected(imported_value.clone(), "signed", serde_json::Value::Null); 750 751 let mut claimed = planned.clone(); 752 let active = claim_for(&claimed, 1, 1, 11); 753 claimed.set_signing_claim(active, 11).unwrap(); 754 let claimed_value = serde_json::to_value(&claimed).unwrap(); 755 assert_artifact_json_rejected(claimed_value.clone(), "revision", serde_json::json!(9)); 756 assert_artifact_json_rejected( 757 claimed_value.clone(), 758 "updated_at_unix_ms", 759 serde_json::json!(12), 760 ); 761 assert_artifact_json_rejected( 762 claimed_value.clone(), 763 "signing_state", 764 serde_json::json!("cancelled"), 765 ); 766 767 let mut signing_retry = planned.clone(); 768 let signing_failure = failure(WorkPhase::Signing, FailureClass::Retryable, Some(20)); 769 signing_retry 770 .record_signing_failure( 771 signing_failure.clone(), 772 Some(RetrySchedule::new(NonZeroU32::MIN, 20, signing_failure).unwrap()), 773 11, 774 ) 775 .unwrap(); 776 let retry_value = serde_json::to_value(&signing_retry).unwrap(); 777 let mut wrong_retry_phase = retry_value.clone(); 778 wrong_retry_phase["signing_retry"]["failure"]["phase"] = serde_json::json!("admission"); 779 assert!(serde_json::from_value::<AuthoredArtifact>(wrong_retry_phase).is_err()); 780 assert_artifact_json_rejected(retry_value.clone(), "last_failure", serde_json::Value::Null); 781 782 let mut signed_pending = planned.clone(); 783 signed_pending 784 .record_signed(signed(&authored_plan), 11) 785 .unwrap(); 786 let signed_value = serde_json::to_value(&signed_pending).unwrap(); 787 let mut admission_claimed = signed_pending.clone(); 788 let active = claim_for(&admission_claimed, 2, 2, 12); 789 admission_claimed.set_admission_claim(active, 12).unwrap(); 790 let admission_claimed_value = serde_json::to_value(&admission_claimed).unwrap(); 791 assert_artifact_json_rejected( 792 admission_claimed_value.clone(), 793 "revision", 794 serde_json::json!(9), 795 ); 796 assert_artifact_json_rejected( 797 admission_claimed_value.clone(), 798 "updated_at_unix_ms", 799 serde_json::json!(13), 800 ); 801 assert_artifact_json_rejected( 802 admission_claimed_value.clone(), 803 "admission_state", 804 serde_json::json!("inserted"), 805 ); 806 let mut claim_without_signed = admission_claimed_value; 807 claim_without_signed["signing_state"] = serde_json::json!("planned"); 808 claim_without_signed["signed"] = serde_json::Value::Null; 809 assert!(serde_json::from_value::<AuthoredArtifact>(claim_without_signed).is_err()); 810 811 let admission_failure = failure(WorkPhase::Admission, FailureClass::Retryable, Some(30)); 812 let mut admission_retry = signed_pending.clone(); 813 admission_retry 814 .record_admission( 815 AdmissionState::Retryable, 816 Some(admission_failure.clone()), 817 Some(RetrySchedule::new(NonZeroU32::MIN, 30, admission_failure).unwrap()), 818 12, 819 ) 820 .unwrap(); 821 let mut wrong_admission_phase = serde_json::to_value(admission_retry).unwrap(); 822 wrong_admission_phase["admission_retry"]["failure"]["phase"] = serde_json::json!("signing"); 823 assert!(serde_json::from_value::<AuthoredArtifact>(wrong_admission_phase).is_err()); 824 825 let terminal = failure(WorkPhase::Admission, FailureClass::Terminal, None); 826 let mut unexpected_failure = signed_value; 827 unexpected_failure["last_failure"] = serde_json::to_value(terminal).unwrap(); 828 assert!(serde_json::from_value::<AuthoredArtifact>(unexpected_failure).is_err()); 829 } 830 831 #[test] 832 fn authored_transition_guards_reject_every_independent_stale_or_incoherent_input() { 833 let authored_plan = plan(); 834 let mut signing = 835 AuthoredArtifact::planned(artifact_id(60), operation_id(), 0, &authored_plan, 10).unwrap(); 836 let active = claim_for(&signing, 1, 2, 11); 837 signing.set_signing_claim(active, 11).unwrap(); 838 let lower_generation = claim_for(&signing, 2, 1, 21); 839 assert_eq!( 840 signing.set_signing_claim(lower_generation, 21), 841 Err(Error::InvalidAuthoredTransition) 842 ); 843 let wrong_revision = WorkClaim::new( 844 [3; 16], 845 "worker", 846 NonZeroU64::new(3).unwrap(), 847 21, 848 30, 849 NonZeroU64::MIN, 850 ) 851 .unwrap(); 852 assert_eq!( 853 signing.set_signing_claim(wrong_revision, 21), 854 Err(Error::InvalidAuthoredTransition) 855 ); 856 let wrong_time = WorkClaim::new( 857 [4; 16], 858 "worker", 859 NonZeroU64::new(4).unwrap(), 860 22, 861 30, 862 signing.revision(), 863 ) 864 .unwrap(); 865 assert_eq!( 866 signing.set_signing_claim(wrong_time, 21), 867 Err(Error::InvalidAuthoredTransition) 868 ); 869 870 let mut admission = signed_artifact(61); 871 let active = claim_for(&admission, 5, 2, 12); 872 admission.set_admission_claim(active, 12).unwrap(); 873 let lower_generation = claim_for(&admission, 6, 1, 22); 874 assert_eq!( 875 admission.set_admission_claim(lower_generation, 22), 876 Err(Error::InvalidAuthoredTransition) 877 ); 878 let wrong_revision = WorkClaim::new( 879 [7; 16], 880 "worker", 881 NonZeroU64::new(3).unwrap(), 882 22, 883 30, 884 NonZeroU64::MIN, 885 ) 886 .unwrap(); 887 assert_eq!( 888 admission.set_admission_claim(wrong_revision, 22), 889 Err(Error::InvalidAuthoredTransition) 890 ); 891 let wrong_time = WorkClaim::new( 892 [8; 16], 893 "worker", 894 NonZeroU64::new(4).unwrap(), 895 23, 896 30, 897 admission.revision(), 898 ) 899 .unwrap(); 900 assert_eq!( 901 admission.set_admission_claim(wrong_time, 22), 902 Err(Error::InvalidAuthoredTransition) 903 ); 904 905 let invalid_retry_failure = failure(WorkPhase::Admission, FailureClass::Retryable, Some(30)); 906 let other_retry_failure = WorkFailure::new( 907 "other_failure", 908 WorkPhase::Admission, 909 FailureClass::Retryable, 910 Some(30), 911 None, 912 ) 913 .unwrap(); 914 let other_retry = RetrySchedule::new(NonZeroU32::MIN, 30, other_retry_failure).unwrap(); 915 for (state, failure, retry) in [ 916 ( 917 AdmissionState::Retryable, 918 Some(invalid_retry_failure.clone()), 919 Some(other_retry), 920 ), 921 ( 922 AdmissionState::Retryable, 923 Some(failure(WorkPhase::Admission, FailureClass::Terminal, None)), 924 Some(RetrySchedule::new(NonZeroU32::MIN, 30, invalid_retry_failure.clone()).unwrap()), 925 ), 926 ( 927 AdmissionState::Duplicate, 928 Some(failure(WorkPhase::Admission, FailureClass::Terminal, None)), 929 None, 930 ), 931 ( 932 AdmissionState::Cancelled, 933 Some(failure(WorkPhase::Admission, FailureClass::Terminal, None)), 934 Some(RetrySchedule::new(NonZeroU32::MIN, 30, invalid_retry_failure.clone()).unwrap()), 935 ), 936 ] { 937 let mut artifact = signed_artifact(70); 938 assert_eq!( 939 artifact.record_admission(state, failure, retry, 12), 940 Err(Error::InvalidAuthoredTransition) 941 ); 942 } 943 } 944 945 fn delivery_request() -> DeliveryRequest { 946 let targets = TargetSet::new(vec![ 947 Target::nostr_relay("wss://one.example").unwrap(), 948 Target::nostr_relay("wss://two.example").unwrap(), 949 ]) 950 .unwrap(); 951 DeliveryRequest::new( 952 "settlement-delivery", 953 DeliveryPayload::new(signed(&plan())), 954 targets, 955 SatisfactionPolicy::new(SatisfactionClass::Accepted, TargetPolicy::all()), 956 100, 957 ) 958 .unwrap() 959 } 960 961 fn settlement_plan(value: u8, artifact: AuthoredArtifactId) -> AuthoredDeliveryPlan { 962 AuthoredDeliveryPlan::new_bound( 963 AuthoredDeliveryPlanId::new([value; 16]).unwrap(), 964 artifact, 965 delivery_request(), 966 10, 967 ) 968 .unwrap() 969 } 970 971 fn claim_delivery(plan: &AuthoredDeliveryPlan, token: u8) -> WorkClaim { 972 WorkClaim::new( 973 [token; 16], 974 "delivery-worker", 975 NonZeroU64::new(u64::from(token)).unwrap(), 976 11, 977 20, 978 plan.revision(), 979 ) 980 .unwrap() 981 } 982 983 fn delivery_receipt(plan: &AuthoredDeliveryPlan, outcome: DeliveryOutcome) -> DeliveryReceipt { 984 let request = plan.request().unwrap(); 985 DeliveryReceipt::for_request( 986 request, 987 request 988 .target_set() 989 .targets() 990 .iter() 991 .cloned() 992 .map(|target| DeliveryTargetReceipt::attempted(target, outcome.clone())) 993 .collect(), 994 ) 995 .unwrap() 996 } 997 998 #[test] 999 fn settlement_counts_every_artifact_and_delivery_terminal_class() { 1000 let authored_plan = plan(); 1001 let mut artifacts = Vec::new(); 1002 for ordinal in 0_u16..11 { 1003 artifacts.push( 1004 AuthoredArtifact::planned( 1005 artifact_id(80 + ordinal as u8), 1006 operation_id(), 1007 ordinal, 1008 &authored_plan, 1009 10, 1010 ) 1011 .unwrap(), 1012 ); 1013 } 1014 let signing_retry = failure(WorkPhase::Signing, FailureClass::Retryable, Some(20)); 1015 artifacts[1] 1016 .record_signing_failure( 1017 signing_retry.clone(), 1018 Some(RetrySchedule::new(NonZeroU32::MIN, 20, signing_retry).unwrap()), 1019 11, 1020 ) 1021 .unwrap(); 1022 artifacts[2] 1023 .record_signing_failure( 1024 failure(WorkPhase::Signing, FailureClass::Indeterminate, None), 1025 None, 1026 11, 1027 ) 1028 .unwrap(); 1029 artifacts[3] 1030 .record_signing_failure( 1031 failure(WorkPhase::Signing, FailureClass::Terminal, None), 1032 None, 1033 11, 1034 ) 1035 .unwrap(); 1036 artifacts[4].cancel_signing(11).unwrap(); 1037 for artifact in &mut artifacts[5..] { 1038 artifact.record_signed(signed(&authored_plan), 11).unwrap(); 1039 } 1040 let admission_retry = failure(WorkPhase::Admission, FailureClass::Retryable, Some(20)); 1041 artifacts[6] 1042 .record_admission( 1043 AdmissionState::Retryable, 1044 Some(admission_retry.clone()), 1045 Some(RetrySchedule::new(NonZeroU32::MIN, 20, admission_retry).unwrap()), 1046 12, 1047 ) 1048 .unwrap(); 1049 for (index, state) in [ 1050 (7, AdmissionState::Rejected), 1051 (8, AdmissionState::Cancelled), 1052 ] { 1053 artifacts[index] 1054 .record_admission( 1055 state, 1056 Some(failure(WorkPhase::Admission, FailureClass::Terminal, None)), 1057 None, 1058 12, 1059 ) 1060 .unwrap(); 1061 } 1062 artifacts[9] 1063 .record_admission(AdmissionState::Inserted, None, None, 12) 1064 .unwrap(); 1065 artifacts[10] 1066 .record_admission(AdmissionState::Duplicate, None, None, 12) 1067 .unwrap(); 1068 1069 let operation = AuthoredOperation::new( 1070 operation_id(), 1071 artifacts 1072 .iter() 1073 .map(AuthoredArtifact::artifact_id) 1074 .collect(), 1075 10, 1076 ) 1077 .unwrap(); 1078 let artifact = artifacts[9].artifact_id(); 1079 let mut plans: Vec<_> = (100_u8..106) 1080 .map(|value| settlement_plan(value, artifact)) 1081 .collect(); 1082 let active = claim_delivery(&plans[1], 1); 1083 plans[1].claim(active.clone(), 11).unwrap(); 1084 let pending = delivery_receipt(&plans[1], DeliveryOutcome::unavailable()); 1085 let retry_failure = failure(WorkPhase::Delivery, FailureClass::Retryable, Some(20)); 1086 plans[1] 1087 .apply_receipt( 1088 active.token(), 1089 active.generation(), 1090 active.row_revision(), 1091 pending, 1092 Some(RetrySchedule::new(NonZeroU32::MIN, 20, retry_failure).unwrap()), 1093 12, 1094 ) 1095 .unwrap(); 1096 let active = claim_delivery(&plans[2], 2); 1097 plans[2].claim(active.clone(), 11).unwrap(); 1098 let accepted = delivery_receipt(&plans[2], DeliveryOutcome::accepted()); 1099 plans[2] 1100 .apply_receipt( 1101 active.token(), 1102 active.generation(), 1103 active.row_revision(), 1104 accepted, 1105 None, 1106 12, 1107 ) 1108 .unwrap(); 1109 let active = claim_delivery(&plans[3], 3); 1110 plans[3].claim(active.clone(), 11).unwrap(); 1111 let rejected = delivery_receipt(&plans[3], DeliveryOutcome::rejected()); 1112 plans[3] 1113 .apply_receipt( 1114 active.token(), 1115 active.generation(), 1116 active.row_revision(), 1117 rejected, 1118 None, 1119 12, 1120 ) 1121 .unwrap(); 1122 let active = claim_delivery(&plans[4], 4); 1123 plans[4].claim(active.clone(), 11).unwrap(); 1124 let request = plans[4].request().unwrap(); 1125 let terminal = SinkFailure::for_request( 1126 request, 1127 "terminal", 1128 Retryability::Terminal, 1129 None, 1130 None, 1131 Vec::new(), 1132 ) 1133 .unwrap(); 1134 plans[4] 1135 .apply_sink_failure( 1136 active.token(), 1137 active.generation(), 1138 active.row_revision(), 1139 terminal, 1140 None, 1141 12, 1142 ) 1143 .unwrap(); 1144 plans[5].cancel(11).unwrap(); 1145 1146 let settlement = 1147 OperationSettlement::evaluate_complete(&operation, &artifacts, &plans).unwrap(); 1148 assert_eq!(settlement.artifacts(), 11); 1149 assert_eq!(settlement.signed(), 6); 1150 assert_eq!(settlement.admitted(), 2); 1151 assert_eq!(settlement.pending(), 2); 1152 assert_eq!(settlement.retryable(), 2); 1153 assert_eq!(settlement.indeterminate(), 1); 1154 assert_eq!(settlement.failed_terminal(), 2); 1155 assert_eq!(settlement.cancelled(), 2); 1156 assert_eq!(settlement.delivery_plans(), 6); 1157 assert_eq!(settlement.delivery_satisfied(), 1); 1158 assert_eq!(settlement.delivery_pending(), 1); 1159 assert_eq!(settlement.delivery_retryable(), 1); 1160 assert_eq!(settlement.delivery_exhausted(), 1); 1161 assert_eq!(settlement.delivery_failed_terminal(), 1); 1162 assert_eq!(settlement.delivery_cancelled(), 1); 1163 assert!(!settlement.is_settled()); 1164 assert!(settlement.has_failures()); 1165 assert!(!settlement.is_successful()); 1166 1167 assert_eq!( 1168 OperationSettlement::evaluate_complete( 1169 &operation, 1170 &artifacts, 1171 &[plans[0].clone(), plans[0].clone()] 1172 ), 1173 Err(Error::InvalidAuthoredOperation) 1174 ); 1175 let foreign = settlement_plan(110, artifact_id(120)); 1176 assert_eq!( 1177 OperationSettlement::evaluate_complete(&operation, &artifacts, &[foreign]), 1178 Err(Error::InvalidAuthoredOperation) 1179 ); 1180 }