authored.rs (65048B)
1 //! Durable authored-operation, artifact, claim, failure, and settlement models. 2 3 use core::num::{NonZeroU32, NonZeroU64}; 4 use radroots_event::SignedEvent; 5 use radroots_event_codec::authoring::{AuthoredEventPlan, HistoricalPlanIntegrity, PlanWireV1}; 6 use sha2::{Digest, Sha256}; 7 use std::{collections::BTreeSet, string::String, vec::Vec}; 8 9 use crate::{ 10 Error, 11 authored_delivery::{AuthoredDeliveryPlan, AuthoredDeliveryState}, 12 journal::OperationInstanceId, 13 }; 14 15 pub const AUTHORED_OPERATION_ARTIFACTS_MAX: usize = 256; 16 pub const WORK_CLAIM_OWNER_MAX_BYTES: usize = 128; 17 pub const WORK_FAILURE_CODE_MAX_BYTES: usize = 64; 18 pub const WORK_FAILURE_DIAGNOSTIC_MAX_BYTES: usize = 1_024; 19 20 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] 21 #[cfg_attr(feature = "serde", serde(try_from = "[u8; 16]", into = "[u8; 16]"))] 22 #[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)] 23 pub struct AuthoredArtifactId([u8; 16]); 24 25 impl AuthoredArtifactId { 26 pub const fn new(value: [u8; 16]) -> Result<Self, Error> { 27 if all_zero(&value) { 28 Err(Error::InvalidAuthoredArtifact) 29 } else { 30 Ok(Self(value)) 31 } 32 } 33 34 pub const fn as_bytes(&self) -> &[u8; 16] { 35 &self.0 36 } 37 } 38 39 impl TryFrom<[u8; 16]> for AuthoredArtifactId { 40 type Error = Error; 41 42 fn try_from(value: [u8; 16]) -> Result<Self, Self::Error> { 43 Self::new(value) 44 } 45 } 46 47 impl From<AuthoredArtifactId> for [u8; 16] { 48 fn from(value: AuthoredArtifactId) -> Self { 49 value.0 50 } 51 } 52 53 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] 54 #[derive(Clone, Copy, Debug, Eq, PartialEq)] 55 #[cfg_attr(feature = "serde", serde(rename_all = "snake_case"))] 56 pub enum ArtifactOrigin { 57 Planned, 58 ImportedSigned, 59 } 60 61 impl ArtifactOrigin { 62 pub const fn is_resignable(self) -> bool { 63 matches!(self, Self::Planned) 64 } 65 } 66 67 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] 68 #[derive(Clone, Copy, Debug, Eq, PartialEq)] 69 #[cfg_attr(feature = "serde", serde(rename_all = "snake_case"))] 70 pub enum SigningState { 71 Planned, 72 Signed, 73 Retryable, 74 Indeterminate, 75 FailedTerminal, 76 Cancelled, 77 } 78 79 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] 80 #[derive(Clone, Copy, Debug, Eq, PartialEq)] 81 #[cfg_attr(feature = "serde", serde(rename_all = "snake_case"))] 82 pub enum AdmissionState { 83 Pending, 84 Inserted, 85 Duplicate, 86 Retryable, 87 Rejected, 88 Cancelled, 89 } 90 91 impl AdmissionState { 92 pub const fn is_admitted(self) -> bool { 93 matches!(self, Self::Inserted | Self::Duplicate) 94 } 95 } 96 97 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] 98 #[derive(Clone, Copy, Debug, Eq, PartialEq)] 99 #[cfg_attr(feature = "serde", serde(rename_all = "snake_case"))] 100 pub enum WorkPhase { 101 Signing, 102 Admission, 103 Delivery, 104 } 105 106 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] 107 #[derive(Clone, Copy, Debug, Eq, PartialEq)] 108 #[cfg_attr(feature = "serde", serde(rename_all = "snake_case"))] 109 pub enum FailureClass { 110 Retryable, 111 Terminal, 112 Indeterminate, 113 } 114 115 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] 116 #[cfg_attr( 117 feature = "serde", 118 serde(try_from = "WorkClaimWire", into = "WorkClaimWire") 119 )] 120 #[derive(Clone, Debug, Eq, PartialEq)] 121 pub struct WorkClaim { 122 token: [u8; 16], 123 owner: String, 124 generation: NonZeroU64, 125 acquired_at_unix_ms: u64, 126 expires_at_unix_ms: u64, 127 row_revision: NonZeroU64, 128 } 129 130 impl WorkClaim { 131 pub fn new( 132 token: [u8; 16], 133 owner: impl Into<String>, 134 generation: NonZeroU64, 135 acquired_at_unix_ms: u64, 136 expires_at_unix_ms: u64, 137 row_revision: NonZeroU64, 138 ) -> Result<Self, Error> { 139 let claim = Self { 140 token, 141 owner: owner.into(), 142 generation, 143 acquired_at_unix_ms, 144 expires_at_unix_ms, 145 row_revision, 146 }; 147 claim.validate()?; 148 Ok(claim) 149 } 150 151 pub fn validate(&self) -> Result<(), Error> { 152 if all_zero(&self.token) 153 || !valid_text(self.owner.as_str(), WORK_CLAIM_OWNER_MAX_BYTES) 154 || self.acquired_at_unix_ms == 0 155 || self.expires_at_unix_ms <= self.acquired_at_unix_ms 156 { 157 return Err(Error::InvalidWorkClaim); 158 } 159 Ok(()) 160 } 161 162 pub const fn token(&self) -> &[u8; 16] { 163 &self.token 164 } 165 pub fn owner(&self) -> &str { 166 self.owner.as_str() 167 } 168 pub const fn generation(&self) -> NonZeroU64 { 169 self.generation 170 } 171 pub const fn acquired_at_unix_ms(&self) -> u64 { 172 self.acquired_at_unix_ms 173 } 174 pub const fn expires_at_unix_ms(&self) -> u64 { 175 self.expires_at_unix_ms 176 } 177 pub const fn row_revision(&self) -> NonZeroU64 { 178 self.row_revision 179 } 180 181 pub fn matches_fence( 182 &self, 183 token: &[u8; 16], 184 generation: NonZeroU64, 185 row_revision: NonZeroU64, 186 now_unix_ms: u64, 187 ) -> bool { 188 &self.token == token 189 && self.generation == generation 190 && self.row_revision == row_revision 191 && now_unix_ms >= self.acquired_at_unix_ms 192 && now_unix_ms < self.expires_at_unix_ms 193 } 194 } 195 196 #[cfg(feature = "serde")] 197 #[derive(serde::Serialize, serde::Deserialize)] 198 struct WorkClaimWire { 199 token: [u8; 16], 200 owner: String, 201 generation: NonZeroU64, 202 acquired_at_unix_ms: u64, 203 expires_at_unix_ms: u64, 204 row_revision: NonZeroU64, 205 } 206 207 #[cfg(feature = "serde")] 208 impl TryFrom<WorkClaimWire> for WorkClaim { 209 type Error = Error; 210 fn try_from(value: WorkClaimWire) -> Result<Self, Self::Error> { 211 Self::new( 212 value.token, 213 value.owner, 214 value.generation, 215 value.acquired_at_unix_ms, 216 value.expires_at_unix_ms, 217 value.row_revision, 218 ) 219 } 220 } 221 222 #[cfg(feature = "serde")] 223 impl From<WorkClaim> for WorkClaimWire { 224 fn from(value: WorkClaim) -> Self { 225 Self { 226 token: value.token, 227 owner: value.owner, 228 generation: value.generation, 229 acquired_at_unix_ms: value.acquired_at_unix_ms, 230 expires_at_unix_ms: value.expires_at_unix_ms, 231 row_revision: value.row_revision, 232 } 233 } 234 } 235 236 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] 237 #[cfg_attr( 238 feature = "serde", 239 serde(try_from = "WorkFailureWire", into = "WorkFailureWire") 240 )] 241 #[derive(Clone, Debug, Eq, PartialEq)] 242 pub struct WorkFailure { 243 code: String, 244 phase: WorkPhase, 245 class: FailureClass, 246 retry_after_unix_ms: Option<u64>, 247 diagnostic: Option<String>, 248 } 249 250 impl WorkFailure { 251 pub fn new( 252 code: impl Into<String>, 253 phase: WorkPhase, 254 class: FailureClass, 255 retry_after_unix_ms: Option<u64>, 256 diagnostic: Option<String>, 257 ) -> Result<Self, Error> { 258 let failure = Self { 259 code: code.into(), 260 phase, 261 class, 262 retry_after_unix_ms, 263 diagnostic, 264 }; 265 failure.validate()?; 266 Ok(failure) 267 } 268 269 pub fn validate(&self) -> Result<(), Error> { 270 if !valid_code(self.code.as_str()) 271 || matches!(self.retry_after_unix_ms, Some(0)) 272 || (self.retry_after_unix_ms.is_some() 273 && !matches!(self.class, FailureClass::Retryable)) 274 || self 275 .diagnostic 276 .as_deref() 277 .is_some_and(|value| !valid_text(value, WORK_FAILURE_DIAGNOSTIC_MAX_BYTES)) 278 { 279 return Err(Error::InvalidWorkFailure); 280 } 281 Ok(()) 282 } 283 284 pub fn code(&self) -> &str { 285 self.code.as_str() 286 } 287 pub const fn phase(&self) -> WorkPhase { 288 self.phase 289 } 290 pub const fn class(&self) -> FailureClass { 291 self.class 292 } 293 pub const fn retry_after_unix_ms(&self) -> Option<u64> { 294 self.retry_after_unix_ms 295 } 296 pub fn diagnostic(&self) -> Option<&str> { 297 self.diagnostic.as_deref() 298 } 299 } 300 301 #[cfg(feature = "serde")] 302 #[derive(serde::Serialize, serde::Deserialize)] 303 struct WorkFailureWire { 304 code: String, 305 phase: WorkPhase, 306 class: FailureClass, 307 retry_after_unix_ms: Option<u64>, 308 diagnostic: Option<String>, 309 } 310 311 #[cfg(feature = "serde")] 312 impl TryFrom<WorkFailureWire> for WorkFailure { 313 type Error = Error; 314 fn try_from(value: WorkFailureWire) -> Result<Self, Self::Error> { 315 Self::new( 316 value.code, 317 value.phase, 318 value.class, 319 value.retry_after_unix_ms, 320 value.diagnostic, 321 ) 322 } 323 } 324 325 #[cfg(feature = "serde")] 326 impl From<WorkFailure> for WorkFailureWire { 327 fn from(value: WorkFailure) -> Self { 328 Self { 329 code: value.code, 330 phase: value.phase, 331 class: value.class, 332 retry_after_unix_ms: value.retry_after_unix_ms, 333 diagnostic: value.diagnostic, 334 } 335 } 336 } 337 338 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] 339 #[cfg_attr( 340 feature = "serde", 341 serde(try_from = "RetryScheduleWire", into = "RetryScheduleWire") 342 )] 343 #[derive(Clone, Debug, Eq, PartialEq)] 344 pub struct RetrySchedule { 345 attempt: NonZeroU32, 346 not_before_unix_ms: u64, 347 failure: WorkFailure, 348 } 349 350 impl RetrySchedule { 351 pub fn new( 352 attempt: NonZeroU32, 353 not_before_unix_ms: u64, 354 failure: WorkFailure, 355 ) -> Result<Self, Error> { 356 if not_before_unix_ms == 0 357 || !matches!(failure.class(), FailureClass::Retryable) 358 || failure 359 .retry_after_unix_ms() 360 .is_some_and(|retry_after| retry_after != not_before_unix_ms) 361 { 362 return Err(Error::InvalidRetrySchedule); 363 } 364 Ok(Self { 365 attempt, 366 not_before_unix_ms, 367 failure, 368 }) 369 } 370 371 pub fn next_attempt( 372 &self, 373 not_before_unix_ms: u64, 374 failure: WorkFailure, 375 ) -> Result<Self, Error> { 376 let attempt = self 377 .attempt 378 .get() 379 .checked_add(1) 380 .and_then(NonZeroU32::new) 381 .ok_or(Error::InvalidRetrySchedule)?; 382 Self::new(attempt, not_before_unix_ms, failure) 383 } 384 385 pub const fn attempt(&self) -> NonZeroU32 { 386 self.attempt 387 } 388 pub const fn not_before_unix_ms(&self) -> u64 { 389 self.not_before_unix_ms 390 } 391 pub const fn failure(&self) -> &WorkFailure { 392 &self.failure 393 } 394 } 395 396 #[cfg(feature = "serde")] 397 #[derive(serde::Serialize, serde::Deserialize)] 398 struct RetryScheduleWire { 399 attempt: NonZeroU32, 400 not_before_unix_ms: u64, 401 failure: WorkFailure, 402 } 403 404 #[cfg(feature = "serde")] 405 impl TryFrom<RetryScheduleWire> for RetrySchedule { 406 type Error = Error; 407 fn try_from(value: RetryScheduleWire) -> Result<Self, Self::Error> { 408 Self::new(value.attempt, value.not_before_unix_ms, value.failure) 409 } 410 } 411 412 #[cfg(feature = "serde")] 413 impl From<RetrySchedule> for RetryScheduleWire { 414 fn from(value: RetrySchedule) -> Self { 415 Self { 416 attempt: value.attempt, 417 not_before_unix_ms: value.not_before_unix_ms, 418 failure: value.failure, 419 } 420 } 421 } 422 423 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] 424 #[cfg_attr( 425 feature = "serde", 426 serde(try_from = "ExactSignedArtifactWire", into = "ExactSignedArtifactWire") 427 )] 428 #[derive(Clone, Debug, Eq, PartialEq)] 429 pub struct ExactSignedArtifact { 430 event: SignedEvent, 431 raw_json_sha256: [u8; 32], 432 } 433 434 #[cfg(feature = "serde")] 435 #[derive(serde::Serialize, serde::Deserialize)] 436 struct ExactSignedArtifactWire { 437 event: SignedEvent, 438 raw_json_sha256: [u8; 32], 439 } 440 441 #[cfg(feature = "serde")] 442 impl TryFrom<ExactSignedArtifactWire> for ExactSignedArtifact { 443 type Error = Error; 444 fn try_from(value: ExactSignedArtifactWire) -> Result<Self, Self::Error> { 445 Self::reconstruct(value.event, value.raw_json_sha256) 446 } 447 } 448 449 #[cfg(feature = "serde")] 450 impl From<ExactSignedArtifact> for ExactSignedArtifactWire { 451 fn from(value: ExactSignedArtifact) -> Self { 452 Self { 453 event: value.event, 454 raw_json_sha256: value.raw_json_sha256, 455 } 456 } 457 } 458 459 /// Exact versioned authored-plan wire retained for durable reconstruction. 460 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] 461 #[cfg_attr( 462 feature = "serde", 463 serde(try_from = "DurableAuthoredPlanWire", into = "DurableAuthoredPlanWire") 464 )] 465 #[derive(Clone, Debug, Eq, PartialEq)] 466 pub struct DurableAuthoredPlan { 467 wire_json: Vec<u8>, 468 } 469 470 #[cfg(feature = "serde")] 471 #[derive(serde::Serialize, serde::Deserialize)] 472 struct DurableAuthoredPlanWire { 473 wire_json: Vec<u8>, 474 } 475 476 #[cfg(feature = "serde")] 477 impl TryFrom<DurableAuthoredPlanWire> for DurableAuthoredPlan { 478 type Error = Error; 479 fn try_from(value: DurableAuthoredPlanWire) -> Result<Self, Self::Error> { 480 Self::reconstruct(value.wire_json) 481 } 482 } 483 484 #[cfg(feature = "serde")] 485 impl From<DurableAuthoredPlan> for DurableAuthoredPlanWire { 486 fn from(value: DurableAuthoredPlan) -> Self { 487 Self { 488 wire_json: value.wire_json, 489 } 490 } 491 } 492 493 impl DurableAuthoredPlan { 494 pub fn from_plan(plan: &AuthoredEventPlan) -> Result<Self, Error> { 495 let wire_json = PlanWireV1::from_plan(plan) 496 .to_json() 497 .map_err(|_| Error::InvalidAuthoredArtifact)?; 498 Ok(Self { wire_json }) 499 } 500 501 pub fn reconstruct(wire_json: Vec<u8>) -> Result<Self, Error> { 502 PlanWireV1::from_json(wire_json.as_slice()).map_err(|_| Error::InvalidAuthoredArtifact)?; 503 Ok(Self { wire_json }) 504 } 505 506 pub fn decode(&self) -> Result<HistoricalPlanIntegrity, Error> { 507 PlanWireV1::from_json(self.wire_json.as_slice()).map_err(|_| Error::InvalidAuthoredArtifact) 508 } 509 510 pub fn wire_json(&self) -> &[u8] { 511 self.wire_json.as_slice() 512 } 513 } 514 515 impl ExactSignedArtifact { 516 pub fn new(event: SignedEvent) -> Self { 517 let raw_json_sha256 = Sha256::digest(event.raw_json().as_bytes()).into(); 518 Self { 519 event, 520 raw_json_sha256, 521 } 522 } 523 524 pub fn reconstruct(event: SignedEvent, raw_json_sha256: [u8; 32]) -> Result<Self, Error> { 525 let artifact = Self { 526 event, 527 raw_json_sha256, 528 }; 529 let computed_sha256: [u8; 32] = Sha256::digest(artifact.event.raw_json().as_bytes()).into(); 530 if computed_sha256 != artifact.raw_json_sha256 { 531 return Err(Error::InvalidAuthoredArtifact); 532 } 533 Ok(artifact) 534 } 535 536 pub const fn event(&self) -> &SignedEvent { 537 &self.event 538 } 539 pub const fn raw_json_sha256(&self) -> &[u8; 32] { 540 &self.raw_json_sha256 541 } 542 } 543 544 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] 545 #[cfg_attr( 546 feature = "serde", 547 serde(try_from = "AuthoredArtifactWire", into = "AuthoredArtifactWire") 548 )] 549 #[derive(Clone, Debug, Eq, PartialEq)] 550 pub struct AuthoredArtifact { 551 artifact_id: AuthoredArtifactId, 552 operation_id: OperationInstanceId, 553 ordinal: u16, 554 origin: ArtifactOrigin, 555 plan: Option<DurableAuthoredPlan>, 556 signing_state: SigningState, 557 admission_state: AdmissionState, 558 signed: Option<ExactSignedArtifact>, 559 signing_claim: Option<WorkClaim>, 560 admission_claim: Option<WorkClaim>, 561 signing_retry: Option<RetrySchedule>, 562 admission_retry: Option<RetrySchedule>, 563 last_failure: Option<WorkFailure>, 564 created_at_unix_ms: u64, 565 updated_at_unix_ms: u64, 566 revision: NonZeroU64, 567 } 568 569 #[cfg(feature = "serde")] 570 #[derive(serde::Serialize, serde::Deserialize)] 571 struct AuthoredArtifactWire { 572 artifact_id: AuthoredArtifactId, 573 operation_id: OperationInstanceId, 574 ordinal: u16, 575 origin: ArtifactOrigin, 576 plan: Option<DurableAuthoredPlan>, 577 signing_state: SigningState, 578 admission_state: AdmissionState, 579 signed: Option<ExactSignedArtifact>, 580 signing_claim: Option<WorkClaim>, 581 admission_claim: Option<WorkClaim>, 582 signing_retry: Option<RetrySchedule>, 583 admission_retry: Option<RetrySchedule>, 584 last_failure: Option<WorkFailure>, 585 created_at_unix_ms: u64, 586 updated_at_unix_ms: u64, 587 revision: NonZeroU64, 588 } 589 590 #[cfg(feature = "serde")] 591 impl TryFrom<AuthoredArtifactWire> for AuthoredArtifact { 592 type Error = Error; 593 fn try_from(value: AuthoredArtifactWire) -> Result<Self, Self::Error> { 594 Self::reconstruct(Self { 595 artifact_id: value.artifact_id, 596 operation_id: value.operation_id, 597 ordinal: value.ordinal, 598 origin: value.origin, 599 plan: value.plan, 600 signing_state: value.signing_state, 601 admission_state: value.admission_state, 602 signed: value.signed, 603 signing_claim: value.signing_claim, 604 admission_claim: value.admission_claim, 605 signing_retry: value.signing_retry, 606 admission_retry: value.admission_retry, 607 last_failure: value.last_failure, 608 created_at_unix_ms: value.created_at_unix_ms, 609 updated_at_unix_ms: value.updated_at_unix_ms, 610 revision: value.revision, 611 }) 612 } 613 } 614 615 #[cfg(feature = "serde")] 616 impl From<AuthoredArtifact> for AuthoredArtifactWire { 617 fn from(value: AuthoredArtifact) -> Self { 618 Self { 619 artifact_id: value.artifact_id, 620 operation_id: value.operation_id, 621 ordinal: value.ordinal, 622 origin: value.origin, 623 plan: value.plan, 624 signing_state: value.signing_state, 625 admission_state: value.admission_state, 626 signed: value.signed, 627 signing_claim: value.signing_claim, 628 admission_claim: value.admission_claim, 629 signing_retry: value.signing_retry, 630 admission_retry: value.admission_retry, 631 last_failure: value.last_failure, 632 created_at_unix_ms: value.created_at_unix_ms, 633 updated_at_unix_ms: value.updated_at_unix_ms, 634 revision: value.revision, 635 } 636 } 637 } 638 639 impl AuthoredArtifact { 640 pub fn planned( 641 artifact_id: AuthoredArtifactId, 642 operation_id: OperationInstanceId, 643 ordinal: u16, 644 plan: &AuthoredEventPlan, 645 created_at_unix_ms: u64, 646 ) -> Result<Self, Error> { 647 Self::reconstruct(Self { 648 artifact_id, 649 operation_id, 650 ordinal, 651 origin: ArtifactOrigin::Planned, 652 plan: Some(DurableAuthoredPlan::from_plan(plan)?), 653 signing_state: SigningState::Planned, 654 admission_state: AdmissionState::Pending, 655 signed: None, 656 signing_claim: None, 657 admission_claim: None, 658 signing_retry: None, 659 admission_retry: None, 660 last_failure: None, 661 created_at_unix_ms, 662 updated_at_unix_ms: created_at_unix_ms, 663 revision: NonZeroU64::MIN, 664 }) 665 } 666 667 pub fn imported_signed( 668 artifact_id: AuthoredArtifactId, 669 operation_id: OperationInstanceId, 670 ordinal: u16, 671 event: SignedEvent, 672 created_at_unix_ms: u64, 673 ) -> Result<Self, Error> { 674 Self::reconstruct(Self { 675 artifact_id, 676 operation_id, 677 ordinal, 678 origin: ArtifactOrigin::ImportedSigned, 679 plan: None, 680 signing_state: SigningState::Signed, 681 admission_state: AdmissionState::Pending, 682 signed: Some(ExactSignedArtifact::new(event)), 683 signing_claim: None, 684 admission_claim: None, 685 signing_retry: None, 686 admission_retry: None, 687 last_failure: None, 688 created_at_unix_ms, 689 updated_at_unix_ms: created_at_unix_ms, 690 revision: NonZeroU64::MIN, 691 }) 692 } 693 694 pub fn reconstruct(value: Self) -> Result<Self, Error> { 695 value.validate()?; 696 Ok(value) 697 } 698 699 pub fn validate(&self) -> Result<(), Error> { 700 if self.created_at_unix_ms == 0 || self.updated_at_unix_ms < self.created_at_unix_ms { 701 return Err(Error::InvalidAuthoredArtifact); 702 } 703 match self.origin { 704 ArtifactOrigin::Planned if self.plan.is_none() => { 705 return Err(Error::InvalidAuthoredArtifact); 706 } 707 ArtifactOrigin::ImportedSigned 708 if self.plan.is_some() 709 || self.signing_state != SigningState::Signed 710 || self.signed.is_none() 711 || self.signing_claim.is_some() 712 || self.signing_retry.is_some() => 713 { 714 return Err(Error::InvalidAuthoredArtifact); 715 } 716 ArtifactOrigin::Planned | ArtifactOrigin::ImportedSigned => {} 717 } 718 if self 719 .plan 720 .as_ref() 721 .is_some_and(|plan| plan.decode().is_err()) 722 { 723 return Err(Error::InvalidAuthoredArtifact); 724 } 725 let stopped = matches!( 726 self.signing_state, 727 SigningState::Cancelled | SigningState::FailedTerminal 728 ); 729 if (self.signing_state == SigningState::Signed && self.signed.is_none()) 730 || (self.signed.is_some() && self.signing_state != SigningState::Signed && !stopped) 731 || (stopped && self.admission_state != AdmissionState::Pending) 732 { 733 return Err(Error::InvalidAuthoredArtifact); 734 } 735 if self.signed.is_none() && self.admission_state != AdmissionState::Pending { 736 return Err(Error::InvalidAuthoredArtifact); 737 } 738 if matches!(self.signing_state, SigningState::Retryable) != self.signing_retry.is_some() 739 || matches!(self.admission_state, AdmissionState::Retryable) 740 != self.admission_retry.is_some() 741 { 742 return Err(Error::InvalidAuthoredArtifact); 743 } 744 if self.signing_claim.as_ref().is_some_and(|claim| { 745 claim.validate().is_err() 746 || !matches!( 747 self.signing_state, 748 SigningState::Planned | SigningState::Retryable 749 ) 750 || claim.row_revision().get().checked_add(1) != Some(self.revision.get()) 751 || claim.acquired_at_unix_ms() != self.updated_at_unix_ms 752 }) || self.admission_claim.as_ref().is_some_and(|claim| { 753 claim.validate().is_err() 754 || self.signing_state != SigningState::Signed 755 || self.signed.is_none() 756 || !matches!( 757 self.admission_state, 758 AdmissionState::Pending | AdmissionState::Retryable 759 ) 760 || claim.row_revision().get().checked_add(1) != Some(self.revision.get()) 761 || claim.acquired_at_unix_ms() != self.updated_at_unix_ms 762 }) { 763 return Err(Error::InvalidAuthoredArtifact); 764 } 765 if self 766 .signing_retry 767 .as_ref() 768 .is_some_and(|retry| retry.failure().phase() != WorkPhase::Signing) 769 || self 770 .admission_retry 771 .as_ref() 772 .is_some_and(|retry| retry.failure().phase() != WorkPhase::Admission) 773 { 774 return Err(Error::InvalidAuthoredArtifact); 775 } 776 if self 777 .last_failure 778 .as_ref() 779 .is_some_and(|failure| failure.validate().is_err()) 780 { 781 return Err(Error::InvalidAuthoredArtifact); 782 } 783 let expected_failure = match self.signing_state { 784 SigningState::Retryable => self.signing_retry.as_ref().map(RetrySchedule::failure), 785 SigningState::Indeterminate => self.last_failure.as_ref().filter(|failure| { 786 failure.phase() == WorkPhase::Signing 787 && failure.class() == FailureClass::Indeterminate 788 }), 789 SigningState::FailedTerminal => self.last_failure.as_ref().filter(|failure| { 790 failure.phase() == WorkPhase::Signing && failure.class() == FailureClass::Terminal 791 }), 792 SigningState::Signed => match self.admission_state { 793 AdmissionState::Retryable => { 794 self.admission_retry.as_ref().map(RetrySchedule::failure) 795 } 796 AdmissionState::Rejected | AdmissionState::Cancelled => { 797 self.last_failure.as_ref().filter(|failure| { 798 failure.phase() == WorkPhase::Admission 799 && failure.class() == FailureClass::Terminal 800 }) 801 } 802 AdmissionState::Pending | AdmissionState::Inserted | AdmissionState::Duplicate => { 803 None 804 } 805 }, 806 SigningState::Planned | SigningState::Cancelled => None, 807 }; 808 let failure_required = matches!( 809 self.signing_state, 810 SigningState::Retryable | SigningState::Indeterminate | SigningState::FailedTerminal 811 ) || (self.signing_state == SigningState::Signed 812 && matches!( 813 self.admission_state, 814 AdmissionState::Retryable | AdmissionState::Rejected | AdmissionState::Cancelled 815 )); 816 if failure_required { 817 if expected_failure.is_none() || expected_failure != self.last_failure.as_ref() { 818 return Err(Error::InvalidAuthoredArtifact); 819 } 820 } else if self.last_failure.is_some() { 821 return Err(Error::InvalidAuthoredArtifact); 822 } 823 if let Some(signed) = &self.signed { 824 ExactSignedArtifact::reconstruct(signed.event.clone(), signed.raw_json_sha256)?; 825 if let Some(plan) = &self.plan { 826 let integrity = plan.decode()?; 827 let plan = integrity.plan(); 828 let event = signed.event(); 829 if event.id() != plan.expected_event_id() 830 || event.pubkey() != plan.author() 831 || event.created_at() != plan.created_at() 832 || event.wire().kind != plan.body().kind() 833 || event.wire().tags != plan.body().tags() 834 || event.wire().content != plan.body().content() 835 { 836 return Err(Error::InvalidAuthoredArtifact); 837 } 838 } 839 } 840 Ok(()) 841 } 842 843 pub fn record_signed(&mut self, event: SignedEvent, at_unix_ms: u64) -> Result<(), Error> { 844 if !self.origin.is_resignable() 845 || !matches!( 846 self.signing_state, 847 SigningState::Planned | SigningState::Retryable 848 ) 849 { 850 return Err(Error::InvalidAuthoredTransition); 851 } 852 let previous = self.clone(); 853 self.signed = Some(ExactSignedArtifact::new(event)); 854 self.signing_state = SigningState::Signed; 855 self.signing_claim = None; 856 self.signing_retry = None; 857 self.last_failure = None; 858 if let Err(error) = self.advance(at_unix_ms) { 859 *self = previous; 860 return Err(error); 861 } 862 if let Err(error) = self.validate() { 863 *self = previous; 864 return Err(error); 865 } 866 Ok(()) 867 } 868 869 pub(crate) fn record_signed_fact( 870 &mut self, 871 event: SignedEvent, 872 observed_at_unix_ms: u64, 873 ) -> Result<(), Error> { 874 if !self.origin.is_resignable() { 875 return Err(Error::InvalidAuthoredTransition); 876 } 877 if let Some(existing) = &self.signed { 878 return if existing.event() == &event { 879 Ok(()) 880 } else { 881 Err(Error::AtomicCommitConflict) 882 }; 883 } 884 let previous = self.clone(); 885 self.signed = Some(ExactSignedArtifact::new(event)); 886 self.signing_claim = None; 887 self.signing_retry = None; 888 if !matches!( 889 self.signing_state, 890 SigningState::Cancelled | SigningState::FailedTerminal 891 ) { 892 self.signing_state = SigningState::Signed; 893 self.last_failure = None; 894 } 895 // Out-of-order observation must not move row time backwards or erase facts. 896 let at_unix_ms = observed_at_unix_ms.max(self.updated_at_unix_ms); 897 if let Err(error) = self.advance(at_unix_ms).and_then(|()| self.validate()) { 898 *self = previous; 899 return Err(error); 900 } 901 Ok(()) 902 } 903 904 pub fn set_signing_claim(&mut self, claim: WorkClaim, at_unix_ms: u64) -> Result<(), Error> { 905 let existing_blocks = self.signing_claim.as_ref().is_some_and(|existing| { 906 at_unix_ms < existing.expires_at_unix_ms() 907 || claim.generation() <= existing.generation() 908 }); 909 if !self.origin.is_resignable() 910 || !matches!( 911 self.signing_state, 912 SigningState::Planned | SigningState::Retryable 913 ) 914 || existing_blocks 915 || claim.row_revision() != self.revision 916 || claim.acquired_at_unix_ms() != at_unix_ms 917 || self 918 .signing_retry 919 .as_ref() 920 .is_some_and(|retry| at_unix_ms < retry.not_before_unix_ms()) 921 { 922 return Err(Error::InvalidAuthoredTransition); 923 } 924 claim.validate()?; 925 let previous = self.clone(); 926 self.signing_claim = Some(claim); 927 if let Err(error) = self.advance(at_unix_ms) { 928 *self = previous; 929 return Err(error); 930 } 931 if let Err(error) = self.validate() { 932 *self = previous; 933 return Err(error); 934 } 935 Ok(()) 936 } 937 938 pub fn record_signing_failure( 939 &mut self, 940 failure: WorkFailure, 941 retry: Option<RetrySchedule>, 942 at_unix_ms: u64, 943 ) -> Result<(), Error> { 944 if failure.phase() != WorkPhase::Signing 945 || !self.origin.is_resignable() 946 || !matches!( 947 self.signing_state, 948 SigningState::Planned | SigningState::Retryable 949 ) 950 || retry 951 .as_ref() 952 .is_some_and(|schedule| schedule.failure() != &failure) 953 { 954 return Err(Error::InvalidAuthoredTransition); 955 } 956 let previous = self.clone(); 957 self.signing_state = match failure.class() { 958 FailureClass::Retryable if retry.is_some() => SigningState::Retryable, 959 FailureClass::Terminal if retry.is_none() => SigningState::FailedTerminal, 960 FailureClass::Indeterminate if retry.is_none() => SigningState::Indeterminate, 961 FailureClass::Retryable | FailureClass::Terminal | FailureClass::Indeterminate => { 962 return Err(Error::InvalidAuthoredTransition); 963 } 964 }; 965 self.signing_claim = None; 966 self.signing_retry = retry; 967 self.last_failure = Some(failure); 968 if let Err(error) = self.advance(at_unix_ms) { 969 *self = previous; 970 return Err(error); 971 } 972 if let Err(error) = self.validate() { 973 *self = previous; 974 return Err(error); 975 } 976 Ok(()) 977 } 978 979 pub fn cancel_signing(&mut self, at_unix_ms: u64) -> Result<(), Error> { 980 if !self.origin.is_resignable() 981 || !matches!( 982 self.signing_state, 983 SigningState::Planned | SigningState::Retryable 984 ) 985 { 986 return Err(Error::InvalidAuthoredTransition); 987 } 988 let previous = self.clone(); 989 self.signing_state = SigningState::Cancelled; 990 self.signing_claim = None; 991 self.signing_retry = None; 992 self.last_failure = None; 993 if let Err(error) = self.advance(at_unix_ms) { 994 *self = previous; 995 return Err(error); 996 } 997 if let Err(error) = self.validate() { 998 *self = previous; 999 return Err(error); 1000 } 1001 Ok(()) 1002 } 1003 1004 pub fn set_admission_claim(&mut self, claim: WorkClaim, at_unix_ms: u64) -> Result<(), Error> { 1005 let existing_blocks = self.admission_claim.as_ref().is_some_and(|existing| { 1006 at_unix_ms < existing.expires_at_unix_ms() 1007 || claim.generation() <= existing.generation() 1008 }); 1009 if self.signing_state != SigningState::Signed 1010 || !matches!( 1011 self.admission_state, 1012 AdmissionState::Pending | AdmissionState::Retryable 1013 ) 1014 || existing_blocks 1015 || claim.row_revision() != self.revision 1016 || claim.acquired_at_unix_ms() != at_unix_ms 1017 || self 1018 .admission_retry 1019 .as_ref() 1020 .is_some_and(|retry| at_unix_ms < retry.not_before_unix_ms()) 1021 { 1022 return Err(Error::InvalidAuthoredTransition); 1023 } 1024 claim.validate()?; 1025 let previous = self.clone(); 1026 self.admission_claim = Some(claim); 1027 if let Err(error) = self.advance(at_unix_ms) { 1028 *self = previous; 1029 return Err(error); 1030 } 1031 if let Err(error) = self.validate() { 1032 *self = previous; 1033 return Err(error); 1034 } 1035 Ok(()) 1036 } 1037 1038 pub fn record_admission( 1039 &mut self, 1040 state: AdmissionState, 1041 failure: Option<WorkFailure>, 1042 retry: Option<RetrySchedule>, 1043 at_unix_ms: u64, 1044 ) -> Result<(), Error> { 1045 if self.signing_state != SigningState::Signed 1046 || self.signed.is_none() 1047 || !matches!( 1048 self.admission_state, 1049 AdmissionState::Pending | AdmissionState::Retryable 1050 ) 1051 || matches!(state, AdmissionState::Pending) 1052 || failure 1053 .as_ref() 1054 .is_some_and(|value| value.phase() != WorkPhase::Admission) 1055 || (matches!(state, AdmissionState::Retryable) != retry.is_some()) 1056 || retry.as_ref().is_some_and(|schedule| { 1057 failure 1058 .as_ref() 1059 .is_none_or(|value| schedule.failure() != value) 1060 }) 1061 || (matches!(state, AdmissionState::Retryable) 1062 && !matches!( 1063 failure.as_ref().map(WorkFailure::class), 1064 Some(FailureClass::Retryable) 1065 )) 1066 || (matches!(state, AdmissionState::Rejected | AdmissionState::Cancelled) 1067 && (failure.is_none() || retry.is_some())) 1068 || (matches!(state, AdmissionState::Inserted | AdmissionState::Duplicate) 1069 && (failure.is_some() || retry.is_some())) 1070 { 1071 return Err(Error::InvalidAuthoredTransition); 1072 } 1073 let previous = self.clone(); 1074 self.admission_state = state; 1075 self.admission_claim = None; 1076 self.admission_retry = retry; 1077 self.last_failure = failure; 1078 if let Err(error) = self.advance(at_unix_ms) { 1079 *self = previous; 1080 return Err(error); 1081 } 1082 if let Err(error) = self.validate() { 1083 *self = previous; 1084 return Err(error); 1085 } 1086 Ok(()) 1087 } 1088 1089 fn advance(&mut self, at_unix_ms: u64) -> Result<(), Error> { 1090 if at_unix_ms < self.updated_at_unix_ms { 1091 return Err(Error::InvalidAuthoredTransition); 1092 } 1093 self.revision = self 1094 .revision 1095 .get() 1096 .checked_add(1) 1097 .and_then(NonZeroU64::new) 1098 .ok_or(Error::InvalidAuthoredTransition)?; 1099 self.updated_at_unix_ms = at_unix_ms; 1100 Ok(()) 1101 } 1102 1103 pub const fn artifact_id(&self) -> AuthoredArtifactId { 1104 self.artifact_id 1105 } 1106 pub const fn operation_id(&self) -> OperationInstanceId { 1107 self.operation_id 1108 } 1109 pub const fn ordinal(&self) -> u16 { 1110 self.ordinal 1111 } 1112 pub const fn origin(&self) -> ArtifactOrigin { 1113 self.origin 1114 } 1115 pub const fn plan(&self) -> Option<&DurableAuthoredPlan> { 1116 self.plan.as_ref() 1117 } 1118 pub const fn signing_state(&self) -> SigningState { 1119 self.signing_state 1120 } 1121 pub const fn admission_state(&self) -> AdmissionState { 1122 self.admission_state 1123 } 1124 pub const fn signed(&self) -> Option<&ExactSignedArtifact> { 1125 self.signed.as_ref() 1126 } 1127 pub const fn signing_claim(&self) -> Option<&WorkClaim> { 1128 self.signing_claim.as_ref() 1129 } 1130 pub const fn admission_claim(&self) -> Option<&WorkClaim> { 1131 self.admission_claim.as_ref() 1132 } 1133 pub const fn signing_retry(&self) -> Option<&RetrySchedule> { 1134 self.signing_retry.as_ref() 1135 } 1136 pub const fn admission_retry(&self) -> Option<&RetrySchedule> { 1137 self.admission_retry.as_ref() 1138 } 1139 pub const fn last_failure(&self) -> Option<&WorkFailure> { 1140 self.last_failure.as_ref() 1141 } 1142 pub const fn created_at_unix_ms(&self) -> u64 { 1143 self.created_at_unix_ms 1144 } 1145 pub const fn updated_at_unix_ms(&self) -> u64 { 1146 self.updated_at_unix_ms 1147 } 1148 pub const fn revision(&self) -> NonZeroU64 { 1149 self.revision 1150 } 1151 } 1152 1153 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] 1154 #[cfg_attr( 1155 feature = "serde", 1156 serde(try_from = "AuthoredOperationWire", into = "AuthoredOperationWire") 1157 )] 1158 #[derive(Clone, Debug, Eq, PartialEq)] 1159 pub struct AuthoredOperation { 1160 operation_id: OperationInstanceId, 1161 artifact_ids: Vec<AuthoredArtifactId>, 1162 created_at_unix_ms: u64, 1163 updated_at_unix_ms: u64, 1164 revision: NonZeroU64, 1165 } 1166 1167 #[cfg(feature = "serde")] 1168 #[derive(serde::Serialize, serde::Deserialize)] 1169 struct AuthoredOperationWire { 1170 operation_id: OperationInstanceId, 1171 artifact_ids: Vec<AuthoredArtifactId>, 1172 created_at_unix_ms: u64, 1173 updated_at_unix_ms: u64, 1174 revision: NonZeroU64, 1175 } 1176 1177 #[cfg(feature = "serde")] 1178 impl TryFrom<AuthoredOperationWire> for AuthoredOperation { 1179 type Error = Error; 1180 fn try_from(value: AuthoredOperationWire) -> Result<Self, Self::Error> { 1181 Self::reconstruct( 1182 value.operation_id, 1183 value.artifact_ids, 1184 value.created_at_unix_ms, 1185 value.updated_at_unix_ms, 1186 value.revision, 1187 ) 1188 } 1189 } 1190 1191 #[cfg(feature = "serde")] 1192 impl From<AuthoredOperation> for AuthoredOperationWire { 1193 fn from(value: AuthoredOperation) -> Self { 1194 Self { 1195 operation_id: value.operation_id, 1196 artifact_ids: value.artifact_ids, 1197 created_at_unix_ms: value.created_at_unix_ms, 1198 updated_at_unix_ms: value.updated_at_unix_ms, 1199 revision: value.revision, 1200 } 1201 } 1202 } 1203 1204 impl AuthoredOperation { 1205 pub fn new( 1206 operation_id: OperationInstanceId, 1207 artifact_ids: Vec<AuthoredArtifactId>, 1208 created_at_unix_ms: u64, 1209 ) -> Result<Self, Error> { 1210 Self::reconstruct( 1211 operation_id, 1212 artifact_ids, 1213 created_at_unix_ms, 1214 created_at_unix_ms, 1215 NonZeroU64::MIN, 1216 ) 1217 } 1218 1219 pub fn reconstruct( 1220 operation_id: OperationInstanceId, 1221 artifact_ids: Vec<AuthoredArtifactId>, 1222 created_at_unix_ms: u64, 1223 updated_at_unix_ms: u64, 1224 revision: NonZeroU64, 1225 ) -> Result<Self, Error> { 1226 if artifact_ids.is_empty() 1227 || artifact_ids.len() > AUTHORED_OPERATION_ARTIFACTS_MAX 1228 || created_at_unix_ms == 0 1229 || updated_at_unix_ms < created_at_unix_ms 1230 || artifact_ids.iter().collect::<BTreeSet<_>>().len() != artifact_ids.len() 1231 { 1232 return Err(Error::InvalidAuthoredOperation); 1233 } 1234 Ok(Self { 1235 operation_id, 1236 artifact_ids, 1237 created_at_unix_ms, 1238 updated_at_unix_ms, 1239 revision, 1240 }) 1241 } 1242 1243 pub const fn operation_id(&self) -> OperationInstanceId { 1244 self.operation_id 1245 } 1246 pub fn artifact_ids(&self) -> &[AuthoredArtifactId] { 1247 self.artifact_ids.as_slice() 1248 } 1249 pub const fn created_at_unix_ms(&self) -> u64 { 1250 self.created_at_unix_ms 1251 } 1252 pub const fn updated_at_unix_ms(&self) -> u64 { 1253 self.updated_at_unix_ms 1254 } 1255 pub const fn revision(&self) -> NonZeroU64 { 1256 self.revision 1257 } 1258 } 1259 1260 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] 1261 #[derive(Clone, Copy, Debug, Eq, PartialEq)] 1262 pub struct OperationSettlement { 1263 artifacts: u16, 1264 signed: u16, 1265 admitted: u16, 1266 pending: u16, 1267 retryable: u16, 1268 indeterminate: u16, 1269 failed_terminal: u16, 1270 cancelled: u16, 1271 delivery_plans: u16, 1272 delivery_satisfied: u16, 1273 delivery_pending: u16, 1274 delivery_retryable: u16, 1275 delivery_exhausted: u16, 1276 delivery_failed_terminal: u16, 1277 delivery_cancelled: u16, 1278 } 1279 1280 impl OperationSettlement { 1281 pub fn evaluate( 1282 operation: &AuthoredOperation, 1283 artifacts: &[AuthoredArtifact], 1284 ) -> Result<Self, Error> { 1285 Self::evaluate_complete(operation, artifacts, &[]) 1286 } 1287 1288 pub fn evaluate_complete( 1289 operation: &AuthoredOperation, 1290 artifacts: &[AuthoredArtifact], 1291 delivery_plans: &[AuthoredDeliveryPlan], 1292 ) -> Result<Self, Error> { 1293 if artifacts.len() != operation.artifact_ids.len() 1294 || artifacts.iter().enumerate().any(|(ordinal, artifact)| { 1295 artifact.operation_id != operation.operation_id 1296 || artifact.artifact_id != operation.artifact_ids[ordinal] 1297 || usize::from(artifact.ordinal) != ordinal 1298 || artifact.validate().is_err() 1299 }) 1300 { 1301 return Err(Error::InvalidAuthoredOperation); 1302 } 1303 let artifact_ids = artifacts 1304 .iter() 1305 .map(AuthoredArtifact::artifact_id) 1306 .collect::<BTreeSet<_>>(); 1307 if delivery_plans 1308 .iter() 1309 .map(AuthoredDeliveryPlan::plan_id) 1310 .collect::<BTreeSet<_>>() 1311 .len() 1312 != delivery_plans.len() 1313 || delivery_plans 1314 .iter() 1315 .any(|plan| !artifact_ids.contains(&plan.artifact_id()) || plan.validate().is_err()) 1316 { 1317 return Err(Error::InvalidAuthoredOperation); 1318 } 1319 let mut settlement = Self { 1320 artifacts: u16::try_from(artifacts.len()) 1321 .map_err(|_| Error::InvalidAuthoredOperation)?, 1322 signed: 0, 1323 admitted: 0, 1324 pending: 0, 1325 retryable: 0, 1326 indeterminate: 0, 1327 failed_terminal: 0, 1328 cancelled: 0, 1329 delivery_plans: u16::try_from(delivery_plans.len()) 1330 .map_err(|_| Error::InvalidAuthoredOperation)?, 1331 delivery_satisfied: 0, 1332 delivery_pending: 0, 1333 delivery_retryable: 0, 1334 delivery_exhausted: 0, 1335 delivery_failed_terminal: 0, 1336 delivery_cancelled: 0, 1337 }; 1338 for artifact in artifacts { 1339 if artifact.signed.is_some() { 1340 settlement.signed += 1; 1341 } 1342 if artifact.admission_state.is_admitted() { 1343 settlement.admitted += 1; 1344 } 1345 match artifact.signing_state { 1346 SigningState::Planned => settlement.pending += 1, 1347 SigningState::Retryable => settlement.retryable += 1, 1348 SigningState::Indeterminate => settlement.indeterminate += 1, 1349 SigningState::FailedTerminal => settlement.failed_terminal += 1, 1350 SigningState::Cancelled => settlement.cancelled += 1, 1351 SigningState::Signed => match artifact.admission_state { 1352 AdmissionState::Pending => settlement.pending += 1, 1353 AdmissionState::Retryable => settlement.retryable += 1, 1354 AdmissionState::Rejected => settlement.failed_terminal += 1, 1355 AdmissionState::Cancelled => settlement.cancelled += 1, 1356 AdmissionState::Inserted | AdmissionState::Duplicate => {} 1357 }, 1358 } 1359 } 1360 for plan in delivery_plans { 1361 match plan.state() { 1362 AuthoredDeliveryState::Pending => settlement.delivery_pending += 1, 1363 AuthoredDeliveryState::Retryable => settlement.delivery_retryable += 1, 1364 AuthoredDeliveryState::Satisfied => settlement.delivery_satisfied += 1, 1365 AuthoredDeliveryState::Exhausted => settlement.delivery_exhausted += 1, 1366 AuthoredDeliveryState::FailedTerminal => { 1367 settlement.delivery_failed_terminal += 1; 1368 } 1369 AuthoredDeliveryState::Cancelled => settlement.delivery_cancelled += 1, 1370 } 1371 } 1372 Ok(settlement) 1373 } 1374 1375 pub const fn artifacts(self) -> u16 { 1376 self.artifacts 1377 } 1378 pub const fn signed(self) -> u16 { 1379 self.signed 1380 } 1381 pub const fn admitted(self) -> u16 { 1382 self.admitted 1383 } 1384 pub const fn pending(self) -> u16 { 1385 self.pending 1386 } 1387 pub const fn retryable(self) -> u16 { 1388 self.retryable 1389 } 1390 pub const fn indeterminate(self) -> u16 { 1391 self.indeterminate 1392 } 1393 pub const fn failed_terminal(self) -> u16 { 1394 self.failed_terminal 1395 } 1396 pub const fn cancelled(self) -> u16 { 1397 self.cancelled 1398 } 1399 pub const fn delivery_plans(self) -> u16 { 1400 self.delivery_plans 1401 } 1402 pub const fn delivery_satisfied(self) -> u16 { 1403 self.delivery_satisfied 1404 } 1405 pub const fn delivery_pending(self) -> u16 { 1406 self.delivery_pending 1407 } 1408 pub const fn delivery_retryable(self) -> u16 { 1409 self.delivery_retryable 1410 } 1411 pub const fn delivery_exhausted(self) -> u16 { 1412 self.delivery_exhausted 1413 } 1414 pub const fn delivery_failed_terminal(self) -> u16 { 1415 self.delivery_failed_terminal 1416 } 1417 pub const fn delivery_cancelled(self) -> u16 { 1418 self.delivery_cancelled 1419 } 1420 pub const fn is_settled(self) -> bool { 1421 self.pending == 0 1422 && self.retryable == 0 1423 && self.indeterminate == 0 1424 && self.delivery_pending == 0 1425 && self.delivery_retryable == 0 1426 } 1427 pub const fn has_failures(self) -> bool { 1428 self.indeterminate != 0 1429 || self.failed_terminal != 0 1430 || self.cancelled != 0 1431 || self.delivery_exhausted != 0 1432 || self.delivery_failed_terminal != 0 1433 || self.delivery_cancelled != 0 1434 } 1435 pub const fn is_successful(self) -> bool { 1436 self.is_settled() 1437 && !self.has_failures() 1438 && self.signed == self.artifacts 1439 && self.admitted == self.artifacts 1440 && self.delivery_satisfied == self.delivery_plans 1441 } 1442 } 1443 1444 const fn all_zero(bytes: &[u8; 16]) -> bool { 1445 let mut index = 0; 1446 while index < bytes.len() { 1447 if bytes[index] != 0 { 1448 return false; 1449 } 1450 index += 1; 1451 } 1452 true 1453 } 1454 1455 fn valid_text(value: &str, max: usize) -> bool { 1456 !value.is_empty() 1457 && value.len() <= max 1458 && value == value.trim() 1459 && !value.chars().any(char::is_control) 1460 } 1461 1462 fn valid_code(value: &str) -> bool { 1463 !value.is_empty() 1464 && value.len() <= WORK_FAILURE_CODE_MAX_BYTES 1465 && value.bytes().all(|byte| { 1466 byte.is_ascii_lowercase() || byte.is_ascii_digit() || matches!(byte, b'_' | b'-' | b'.') 1467 }) 1468 } 1469 1470 #[cfg(test)] 1471 mod invariant_tests { 1472 use super::*; 1473 use radroots_event::{GenericEventDraft, wire::v1::Nip01EventWire}; 1474 1475 const AUTHOR: &str = "585591529da0bab31b3b1b1f986611cf5f435dca84f978c89ee8a40cca7103df"; 1476 1477 fn operation_id() -> OperationInstanceId { 1478 OperationInstanceId::new([1; 16]).unwrap() 1479 } 1480 1481 fn artifact_id(value: u8) -> AuthoredArtifactId { 1482 AuthoredArtifactId::new([value; 16]).unwrap() 1483 } 1484 1485 fn plan() -> AuthoredEventPlan { 1486 AuthoredEventPlan::from_generic( 1487 GenericEventDraft::new( 1488 "radroots.social.geochat.v1", 1489 20_000, 1490 1_800_000_100, 1491 Vec::new(), 1492 "authored invariant", 1493 AUTHOR, 1494 ) 1495 .unwrap(), 1496 ) 1497 .unwrap() 1498 } 1499 1500 fn signed(plan: &AuthoredEventPlan) -> SignedEvent { 1501 let wire = Nip01EventWire { 1502 id: plan.expected_event_id().to_hex(), 1503 pubkey: plan.author().to_hex(), 1504 created_at: plan.created_at(), 1505 kind: plan.body().kind(), 1506 tags: plan.body().tags().to_vec(), 1507 content: plan.body().content().to_owned(), 1508 sig: "dd".repeat(64), 1509 extra: Default::default(), 1510 }; 1511 let raw = serde_json::to_string(&wire).unwrap(); 1512 SignedEvent::from_wire_verified_id(wire, raw).unwrap() 1513 } 1514 1515 fn planned(value: u8) -> AuthoredArtifact { 1516 AuthoredArtifact::planned(artifact_id(value), operation_id(), 0, &plan(), 10).unwrap() 1517 } 1518 1519 fn imported(value: u8) -> AuthoredArtifact { 1520 let plan = plan(); 1521 AuthoredArtifact::imported_signed(artifact_id(value), operation_id(), 0, signed(&plan), 10) 1522 .unwrap() 1523 } 1524 1525 fn valid_claim(revision: NonZeroU64, at: u64) -> WorkClaim { 1526 WorkClaim::new([7; 16], "worker", NonZeroU64::MIN, at, at + 10, revision).unwrap() 1527 } 1528 1529 fn failure(phase: WorkPhase, class: FailureClass) -> WorkFailure { 1530 let retry_at = (class == FailureClass::Retryable).then_some(30); 1531 WorkFailure::new("failure", phase, class, retry_at, None).unwrap() 1532 } 1533 1534 fn retry(phase: WorkPhase) -> RetrySchedule { 1535 RetrySchedule::new(NonZeroU32::MIN, 30, failure(phase, FailureClass::Retryable)).unwrap() 1536 } 1537 1538 #[test] 1539 fn validator_rejects_corruption_at_each_nested_invariant() { 1540 let mut value = imported(1); 1541 value.signing_claim = Some(valid_claim(value.revision, 10)); 1542 assert_eq!(value.validate(), Err(Error::InvalidAuthoredArtifact)); 1543 1544 let mut value = imported(2); 1545 value.signing_retry = Some(retry(WorkPhase::Signing)); 1546 assert_eq!(value.validate(), Err(Error::InvalidAuthoredArtifact)); 1547 1548 let mut value = planned(3); 1549 value.plan = Some(DurableAuthoredPlan { 1550 wire_json: vec![0xff], 1551 }); 1552 assert_eq!(value.validate(), Err(Error::InvalidAuthoredArtifact)); 1553 1554 let mut value = imported(4); 1555 value.admission_state = AdmissionState::Retryable; 1556 assert_eq!(value.validate(), Err(Error::InvalidAuthoredArtifact)); 1557 1558 let mut claimed = planned(5); 1559 claimed 1560 .set_signing_claim(valid_claim(claimed.revision, 11), 11) 1561 .unwrap(); 1562 let mut value = claimed.clone(); 1563 value.signing_claim.as_mut().unwrap().token = [0; 16]; 1564 assert_eq!(value.validate(), Err(Error::InvalidAuthoredArtifact)); 1565 let mut value = claimed.clone(); 1566 value.signing_state = SigningState::Cancelled; 1567 assert_eq!(value.validate(), Err(Error::InvalidAuthoredArtifact)); 1568 let mut value = claimed.clone(); 1569 value.signing_claim.as_mut().unwrap().row_revision = value.revision; 1570 assert_eq!(value.validate(), Err(Error::InvalidAuthoredArtifact)); 1571 let mut value = claimed; 1572 value.signing_claim.as_mut().unwrap().acquired_at_unix_ms += 1; 1573 assert_eq!(value.validate(), Err(Error::InvalidAuthoredArtifact)); 1574 1575 let authored_plan = plan(); 1576 let mut admission = planned(6); 1577 admission.record_signed(signed(&authored_plan), 11).unwrap(); 1578 admission 1579 .set_admission_claim(valid_claim(admission.revision, 12), 12) 1580 .unwrap(); 1581 let mut value = admission.clone(); 1582 value.admission_claim.as_mut().unwrap().token = [0; 16]; 1583 assert_eq!(value.validate(), Err(Error::InvalidAuthoredArtifact)); 1584 let mut value = admission.clone(); 1585 value.signed = None; 1586 value.signing_state = SigningState::Planned; 1587 assert_eq!(value.validate(), Err(Error::InvalidAuthoredArtifact)); 1588 let mut value = admission.clone(); 1589 value.admission_state = AdmissionState::Inserted; 1590 assert_eq!(value.validate(), Err(Error::InvalidAuthoredArtifact)); 1591 let mut value = admission.clone(); 1592 value.admission_claim.as_mut().unwrap().row_revision = value.revision; 1593 assert_eq!(value.validate(), Err(Error::InvalidAuthoredArtifact)); 1594 let mut value = admission; 1595 value.admission_claim.as_mut().unwrap().acquired_at_unix_ms += 1; 1596 assert_eq!(value.validate(), Err(Error::InvalidAuthoredArtifact)); 1597 1598 let mut value = planned(7); 1599 value.signing_state = SigningState::Retryable; 1600 value.signing_retry = Some(retry(WorkPhase::Admission)); 1601 value.last_failure = value 1602 .signing_retry 1603 .as_ref() 1604 .map(|schedule| schedule.failure.clone()); 1605 assert_eq!(value.validate(), Err(Error::InvalidAuthoredArtifact)); 1606 1607 let mut value = imported(8); 1608 value.admission_state = AdmissionState::Retryable; 1609 value.admission_retry = Some(retry(WorkPhase::Signing)); 1610 value.last_failure = value 1611 .admission_retry 1612 .as_ref() 1613 .map(|schedule| schedule.failure.clone()); 1614 assert_eq!(value.validate(), Err(Error::InvalidAuthoredArtifact)); 1615 1616 let mut value = planned(9); 1617 value.last_failure = Some(WorkFailure { 1618 code: "INVALID".to_owned(), 1619 phase: WorkPhase::Signing, 1620 class: FailureClass::Terminal, 1621 retry_after_unix_ms: None, 1622 diagnostic: None, 1623 }); 1624 assert_eq!(value.validate(), Err(Error::InvalidAuthoredArtifact)); 1625 1626 let mut value = planned(10); 1627 value.signing_state = SigningState::Indeterminate; 1628 value.last_failure = Some(failure(WorkPhase::Signing, FailureClass::Terminal)); 1629 assert_eq!(value.validate(), Err(Error::InvalidAuthoredArtifact)); 1630 let mut value = planned(11); 1631 value.signing_state = SigningState::FailedTerminal; 1632 value.last_failure = Some(failure(WorkPhase::Signing, FailureClass::Indeterminate)); 1633 assert_eq!(value.validate(), Err(Error::InvalidAuthoredArtifact)); 1634 let mut value = imported(12); 1635 value.admission_state = AdmissionState::Rejected; 1636 value.last_failure = Some(failure(WorkPhase::Signing, FailureClass::Terminal)); 1637 assert_eq!(value.validate(), Err(Error::InvalidAuthoredArtifact)); 1638 let mut value = planned(13); 1639 value.last_failure = Some(failure(WorkPhase::Signing, FailureClass::Terminal)); 1640 assert_eq!(value.validate(), Err(Error::InvalidAuthoredArtifact)); 1641 } 1642 1643 #[test] 1644 fn every_mutating_transition_rolls_back_on_revision_overflow() { 1645 let authored_plan = plan(); 1646 1647 let mut value = planned(20); 1648 value.revision = NonZeroU64::MAX; 1649 let before = value.clone(); 1650 assert_eq!( 1651 value.record_signed(signed(&authored_plan), 11), 1652 Err(Error::InvalidAuthoredTransition) 1653 ); 1654 assert_eq!(value, before); 1655 1656 let mut value = planned(21); 1657 value.revision = NonZeroU64::MAX; 1658 let before = value.clone(); 1659 assert_eq!( 1660 value.set_signing_claim(valid_claim(value.revision, 11), 11), 1661 Err(Error::InvalidAuthoredTransition) 1662 ); 1663 assert_eq!(value, before); 1664 1665 let mut value = planned(22); 1666 value.revision = NonZeroU64::MAX; 1667 let before = value.clone(); 1668 assert_eq!( 1669 value.record_signing_failure( 1670 failure(WorkPhase::Signing, FailureClass::Terminal), 1671 None, 1672 11, 1673 ), 1674 Err(Error::InvalidAuthoredTransition) 1675 ); 1676 assert_eq!(value, before); 1677 1678 let mut value = planned(23); 1679 value.revision = NonZeroU64::MAX; 1680 let before = value.clone(); 1681 assert_eq!( 1682 value.cancel_signing(11), 1683 Err(Error::InvalidAuthoredTransition) 1684 ); 1685 assert_eq!(value, before); 1686 1687 let mut value = imported(24); 1688 value.revision = NonZeroU64::MAX; 1689 let before = value.clone(); 1690 assert_eq!( 1691 value.set_admission_claim(valid_claim(value.revision, 11), 11), 1692 Err(Error::InvalidAuthoredTransition) 1693 ); 1694 assert_eq!(value, before); 1695 1696 let mut value = imported(25); 1697 value.revision = NonZeroU64::MAX; 1698 let before = value.clone(); 1699 assert_eq!( 1700 value.record_admission(AdmissionState::Inserted, None, None, 11), 1701 Err(Error::InvalidAuthoredTransition) 1702 ); 1703 assert_eq!(value, before); 1704 } 1705 1706 #[test] 1707 fn every_mutating_transition_rolls_back_when_post_state_is_invalid() { 1708 let authored_plan = plan(); 1709 let mut values = (30_u8..34).map(planned).collect::<Vec<_>>(); 1710 for value in &mut values { 1711 value.plan = None; 1712 } 1713 1714 let before = values[0].clone(); 1715 assert_eq!( 1716 values[0].record_signed(signed(&authored_plan), 11), 1717 Err(Error::InvalidAuthoredArtifact) 1718 ); 1719 assert_eq!(values[0], before); 1720 1721 let before = values[1].clone(); 1722 let revision = values[1].revision; 1723 assert_eq!( 1724 values[1].set_signing_claim(valid_claim(revision, 11), 11), 1725 Err(Error::InvalidAuthoredArtifact) 1726 ); 1727 assert_eq!(values[1], before); 1728 1729 let before = values[2].clone(); 1730 assert_eq!( 1731 values[2].record_signing_failure( 1732 failure(WorkPhase::Signing, FailureClass::Terminal), 1733 None, 1734 11, 1735 ), 1736 Err(Error::InvalidAuthoredArtifact) 1737 ); 1738 assert_eq!(values[2], before); 1739 1740 let before = values[3].clone(); 1741 assert_eq!( 1742 values[3].cancel_signing(11), 1743 Err(Error::InvalidAuthoredArtifact) 1744 ); 1745 assert_eq!(values[3], before); 1746 1747 let mut value = imported(34); 1748 value.origin = ArtifactOrigin::Planned; 1749 let before = value.clone(); 1750 assert_eq!( 1751 value.set_admission_claim(valid_claim(value.revision, 11), 11), 1752 Err(Error::InvalidAuthoredArtifact) 1753 ); 1754 assert_eq!(value, before); 1755 1756 let mut value = imported(35); 1757 value.origin = ArtifactOrigin::Planned; 1758 let before = value.clone(); 1759 assert_eq!( 1760 value.record_admission(AdmissionState::Inserted, None, None, 11), 1761 Err(Error::InvalidAuthoredArtifact) 1762 ); 1763 assert_eq!(value, before); 1764 } 1765 1766 fn empty_settlement() -> OperationSettlement { 1767 OperationSettlement { 1768 artifacts: 1, 1769 signed: 1, 1770 admitted: 1, 1771 pending: 0, 1772 retryable: 0, 1773 indeterminate: 0, 1774 failed_terminal: 0, 1775 cancelled: 0, 1776 delivery_plans: 0, 1777 delivery_satisfied: 0, 1778 delivery_pending: 0, 1779 delivery_retryable: 0, 1780 delivery_exhausted: 0, 1781 delivery_failed_terminal: 0, 1782 delivery_cancelled: 0, 1783 } 1784 } 1785 1786 #[test] 1787 fn settlement_predicates_evaluate_every_independent_dimension() { 1788 assert!(empty_settlement().is_settled()); 1789 for mutate in [ 1790 |value: &mut OperationSettlement| value.pending = 1, 1791 |value: &mut OperationSettlement| value.retryable = 1, 1792 |value: &mut OperationSettlement| value.indeterminate = 1, 1793 |value: &mut OperationSettlement| value.delivery_pending = 1, 1794 |value: &mut OperationSettlement| value.delivery_retryable = 1, 1795 ] { 1796 let mut value = empty_settlement(); 1797 mutate(&mut value); 1798 assert!(!value.is_settled()); 1799 } 1800 1801 assert!(!empty_settlement().has_failures()); 1802 for mutate in [ 1803 |value: &mut OperationSettlement| value.indeterminate = 1, 1804 |value: &mut OperationSettlement| value.failed_terminal = 1, 1805 |value: &mut OperationSettlement| value.cancelled = 1, 1806 |value: &mut OperationSettlement| value.delivery_exhausted = 1, 1807 |value: &mut OperationSettlement| value.delivery_failed_terminal = 1, 1808 |value: &mut OperationSettlement| value.delivery_cancelled = 1, 1809 ] { 1810 let mut value = empty_settlement(); 1811 mutate(&mut value); 1812 assert!(value.has_failures()); 1813 } 1814 1815 assert!(empty_settlement().is_successful()); 1816 for mutate in [ 1817 |value: &mut OperationSettlement| value.pending = 1, 1818 |value: &mut OperationSettlement| value.failed_terminal = 1, 1819 |value: &mut OperationSettlement| value.signed = 0, 1820 |value: &mut OperationSettlement| value.admitted = 0, 1821 |value: &mut OperationSettlement| value.delivery_plans = 1, 1822 ] { 1823 let mut value = empty_settlement(); 1824 mutate(&mut value); 1825 assert!(!value.is_successful()); 1826 } 1827 } 1828 1829 #[test] 1830 fn settlement_rejects_each_artifact_identity_and_integrity_mismatch() { 1831 let artifact = planned(40); 1832 let operation = 1833 AuthoredOperation::new(operation_id(), vec![artifact.artifact_id()], 10).unwrap(); 1834 assert_eq!( 1835 OperationSettlement::evaluate(&operation, &[]), 1836 Err(Error::InvalidAuthoredOperation) 1837 ); 1838 1839 let mut wrong_operation = artifact.clone(); 1840 wrong_operation.operation_id = OperationInstanceId::new([2; 16]).unwrap(); 1841 assert_eq!( 1842 OperationSettlement::evaluate(&operation, &[wrong_operation]), 1843 Err(Error::InvalidAuthoredOperation) 1844 ); 1845 let mut wrong_id = artifact.clone(); 1846 wrong_id.artifact_id = artifact_id(41); 1847 assert_eq!( 1848 OperationSettlement::evaluate(&operation, &[wrong_id]), 1849 Err(Error::InvalidAuthoredOperation) 1850 ); 1851 let mut wrong_ordinal = artifact.clone(); 1852 wrong_ordinal.ordinal = 1; 1853 assert_eq!( 1854 OperationSettlement::evaluate(&operation, &[wrong_ordinal]), 1855 Err(Error::InvalidAuthoredOperation) 1856 ); 1857 let mut invalid = artifact; 1858 invalid.plan = None; 1859 assert_eq!( 1860 OperationSettlement::evaluate(&operation, &[invalid]), 1861 Err(Error::InvalidAuthoredOperation) 1862 ); 1863 } 1864 1865 #[test] 1866 fn transition_preconditions_distinguish_origin_state_and_signed_presence() { 1867 let authored_plan = plan(); 1868 1869 let mut value = planned(50); 1870 value.signing_state = SigningState::Cancelled; 1871 assert_eq!( 1872 value.record_signed(signed(&authored_plan), 11), 1873 Err(Error::InvalidAuthoredTransition) 1874 ); 1875 1876 let mut value = imported(51); 1877 assert_eq!( 1878 value.set_signing_claim(valid_claim(value.revision, 11), 11), 1879 Err(Error::InvalidAuthoredTransition) 1880 ); 1881 let mut value = planned(52); 1882 value.signing_state = SigningState::Cancelled; 1883 assert_eq!( 1884 value.set_signing_claim(valid_claim(value.revision, 11), 11), 1885 Err(Error::InvalidAuthoredTransition) 1886 ); 1887 1888 let mut value = imported(53); 1889 assert_eq!( 1890 value.record_signing_failure( 1891 failure(WorkPhase::Signing, FailureClass::Terminal), 1892 None, 1893 11, 1894 ), 1895 Err(Error::InvalidAuthoredTransition) 1896 ); 1897 let mut value = planned(54); 1898 value.signing_state = SigningState::Cancelled; 1899 assert_eq!( 1900 value.record_signing_failure( 1901 failure(WorkPhase::Signing, FailureClass::Terminal), 1902 None, 1903 11, 1904 ), 1905 Err(Error::InvalidAuthoredTransition) 1906 ); 1907 1908 let mut value = imported(55); 1909 value.admission_state = AdmissionState::Inserted; 1910 assert_eq!( 1911 value.set_admission_claim(valid_claim(value.revision, 11), 11), 1912 Err(Error::InvalidAuthoredTransition) 1913 ); 1914 1915 let mut missing_signed = imported(56); 1916 missing_signed.signed = None; 1917 assert_eq!( 1918 missing_signed.record_admission(AdmissionState::Inserted, None, None, 11), 1919 Err(Error::InvalidAuthoredTransition) 1920 ); 1921 let mut terminal_admission = imported(57); 1922 terminal_admission.admission_state = AdmissionState::Inserted; 1923 assert_eq!( 1924 terminal_admission.record_admission(AdmissionState::Duplicate, None, None, 11), 1925 Err(Error::InvalidAuthoredTransition) 1926 ); 1927 } 1928 }