lib

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

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 }