authored_atomic.rs (30515B)
1 //! Atomic authored-operation commands with deterministic phase identities. 2 3 mod signing_evidence; 4 pub use signing_evidence::RecordSignedArtifact; 5 mod delivery_evidence; 6 pub use delivery_evidence::RecordDeliveryFact; 7 mod delivery_reconciliation; 8 pub use delivery_reconciliation::ReconcileDeliveryFacts; 9 10 use core::num::NonZeroU64; 11 use radroots_event::SignedEvent; 12 use radroots_transport::BoxFuture; 13 use sha2::{Digest, Sha256}; 14 use std::{collections::BTreeSet, vec::Vec}; 15 16 use crate::{ 17 Error, 18 atomic::{AtomicCommitDigest, AtomicCommitDisposition, AtomicCommitId}, 19 authored::{ 20 AdmissionState, AuthoredArtifact, AuthoredArtifactId, AuthoredOperation, RetrySchedule, 21 WorkFailure, WorkPhase, 22 }, 23 authored_delivery::{AuthoredDeliveryPlan, AuthoredDeliveryPlanId, DeliveryAttemptOutcome}, 24 authored_draft_submission::PrepareFromDraft, 25 journal::OperationInstanceId, 26 }; 27 28 #[derive(Clone, Debug, Eq, PartialEq)] 29 pub struct WorkFence { 30 token: [u8; 16], 31 generation: NonZeroU64, 32 row_revision: NonZeroU64, 33 } 34 35 impl WorkFence { 36 pub const fn new( 37 token: [u8; 16], 38 generation: NonZeroU64, 39 row_revision: NonZeroU64, 40 ) -> Result<Self, Error> { 41 if bytes_are_zero(&token) { 42 return Err(Error::InvalidWorkClaim); 43 } 44 Ok(Self { 45 token, 46 generation, 47 row_revision, 48 }) 49 } 50 pub const fn token(&self) -> &[u8; 16] { 51 &self.token 52 } 53 pub const fn generation(&self) -> NonZeroU64 { 54 self.generation 55 } 56 pub const fn row_revision(&self) -> NonZeroU64 { 57 self.row_revision 58 } 59 } 60 61 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] 62 #[cfg_attr( 63 feature = "serde", 64 serde(try_from = "PrepareWire", into = "PrepareWire") 65 )] 66 #[derive(Clone, Debug, Eq, PartialEq)] 67 pub struct PrepareAuthoredOperation { 68 operation: AuthoredOperation, 69 artifacts: Vec<AuthoredArtifact>, 70 delivery_plans: Vec<AuthoredDeliveryPlan>, 71 input_digest: AtomicCommitDigest, 72 requested_at_unix_ms: u64, 73 } 74 75 #[cfg(feature = "serde")] 76 #[derive(serde::Serialize, serde::Deserialize)] 77 #[serde(deny_unknown_fields)] 78 struct PrepareWire { 79 operation: AuthoredOperation, 80 artifacts: Vec<AuthoredArtifact>, 81 delivery_plans: Vec<AuthoredDeliveryPlan>, 82 input_digest: AtomicCommitDigest, 83 requested_at_unix_ms: u64, 84 } 85 #[cfg(feature = "serde")] 86 impl TryFrom<PrepareWire> for PrepareAuthoredOperation { 87 type Error = Error; 88 fn try_from(v: PrepareWire) -> Result<Self, Error> { 89 Self::new( 90 v.operation, 91 v.artifacts, 92 v.delivery_plans, 93 v.input_digest, 94 v.requested_at_unix_ms, 95 ) 96 } 97 } 98 #[cfg(feature = "serde")] 99 impl From<PrepareAuthoredOperation> for PrepareWire { 100 fn from(v: PrepareAuthoredOperation) -> Self { 101 Self { 102 operation: v.operation, 103 artifacts: v.artifacts, 104 delivery_plans: v.delivery_plans, 105 input_digest: v.input_digest, 106 requested_at_unix_ms: v.requested_at_unix_ms, 107 } 108 } 109 } 110 111 impl PrepareAuthoredOperation { 112 pub fn new( 113 operation: AuthoredOperation, 114 artifacts: Vec<AuthoredArtifact>, 115 delivery_plans: Vec<AuthoredDeliveryPlan>, 116 input_digest: AtomicCommitDigest, 117 requested_at_unix_ms: u64, 118 ) -> Result<Self, Error> { 119 if requested_at_unix_ms == 0 120 || artifacts.len() != operation.artifact_ids().len() 121 || artifacts.iter().enumerate().any(|(ordinal, artifact)| { 122 artifact.operation_id() != operation.operation_id() 123 || artifact.artifact_id() != operation.artifact_ids()[ordinal] 124 || usize::from(artifact.ordinal()) != ordinal 125 || artifact.validate().is_err() 126 }) 127 { 128 return Err(Error::AtomicWorkflowMismatch); 129 } 130 let artifact_ids: BTreeSet<_> = artifacts 131 .iter() 132 .map(AuthoredArtifact::artifact_id) 133 .collect(); 134 let plan_ids: BTreeSet<_> = delivery_plans 135 .iter() 136 .map(AuthoredDeliveryPlan::plan_id) 137 .collect(); 138 if plan_ids.len() != delivery_plans.len() 139 || delivery_plans.iter().any(|plan| { 140 !artifact_ids.contains(&plan.artifact_id()) 141 || plan.validate().is_err() 142 || !plan.delivery_facts().is_empty() 143 }) 144 { 145 return Err(Error::AtomicWorkflowMismatch); 146 } 147 Ok(Self { 148 operation, 149 artifacts, 150 delivery_plans, 151 input_digest, 152 requested_at_unix_ms, 153 }) 154 } 155 pub const fn operation(&self) -> &AuthoredOperation { 156 &self.operation 157 } 158 pub fn artifacts(&self) -> &[AuthoredArtifact] { 159 self.artifacts.as_slice() 160 } 161 pub fn delivery_plans(&self) -> &[AuthoredDeliveryPlan] { 162 self.delivery_plans.as_slice() 163 } 164 pub const fn input_digest(&self) -> AtomicCommitDigest { 165 self.input_digest 166 } 167 pub const fn requested_at_unix_ms(&self) -> u64 { 168 self.requested_at_unix_ms 169 } 170 } 171 172 #[derive(Clone, Debug, Eq, PartialEq)] 173 pub struct ApplySignedArtifact { 174 artifact_id: AuthoredArtifactId, 175 fence: WorkFence, 176 event: SignedEvent, 177 applied_at_unix_ms: u64, 178 } 179 180 impl ApplySignedArtifact { 181 pub fn new( 182 artifact_id: AuthoredArtifactId, 183 fence: WorkFence, 184 event: SignedEvent, 185 applied_at_unix_ms: u64, 186 ) -> Result<Self, Error> { 187 if applied_at_unix_ms == 0 { 188 return Err(Error::AtomicWorkflowMismatch); 189 } 190 Ok(Self { 191 artifact_id, 192 fence, 193 event, 194 applied_at_unix_ms, 195 }) 196 } 197 pub const fn artifact_id(&self) -> AuthoredArtifactId { 198 self.artifact_id 199 } 200 pub const fn fence(&self) -> &WorkFence { 201 &self.fence 202 } 203 pub const fn event(&self) -> &SignedEvent { 204 &self.event 205 } 206 pub const fn applied_at_unix_ms(&self) -> u64 { 207 self.applied_at_unix_ms 208 } 209 } 210 211 #[derive(Clone, Debug, Eq, PartialEq)] 212 pub struct ApplyAdmissionResult { 213 artifact_id: AuthoredArtifactId, 214 fence: WorkFence, 215 state: AdmissionState, 216 failure: Option<WorkFailure>, 217 retry: Option<RetrySchedule>, 218 applied_at_unix_ms: u64, 219 } 220 221 impl ApplyAdmissionResult { 222 pub fn new( 223 artifact_id: AuthoredArtifactId, 224 fence: WorkFence, 225 state: AdmissionState, 226 failure: Option<WorkFailure>, 227 retry: Option<RetrySchedule>, 228 applied_at_unix_ms: u64, 229 ) -> Result<Self, Error> { 230 if applied_at_unix_ms == 0 { 231 return Err(Error::AtomicWorkflowMismatch); 232 } 233 Ok(Self { 234 artifact_id, 235 fence, 236 state, 237 failure, 238 retry, 239 applied_at_unix_ms, 240 }) 241 } 242 pub const fn artifact_id(&self) -> AuthoredArtifactId { 243 self.artifact_id 244 } 245 pub const fn fence(&self) -> &WorkFence { 246 &self.fence 247 } 248 pub const fn state(&self) -> AdmissionState { 249 self.state 250 } 251 pub const fn failure(&self) -> Option<&WorkFailure> { 252 self.failure.as_ref() 253 } 254 pub const fn retry(&self) -> Option<&RetrySchedule> { 255 self.retry.as_ref() 256 } 257 pub const fn applied_at_unix_ms(&self) -> u64 { 258 self.applied_at_unix_ms 259 } 260 } 261 262 #[derive(Clone, Debug, Eq, PartialEq)] 263 pub struct ApplyDeliveryAttempt { 264 plan_id: AuthoredDeliveryPlanId, 265 fence: WorkFence, 266 outcome: DeliveryAttemptOutcome, 267 retry: Option<RetrySchedule>, 268 applied_at_unix_ms: u64, 269 } 270 271 impl ApplyDeliveryAttempt { 272 pub fn new( 273 plan_id: AuthoredDeliveryPlanId, 274 fence: WorkFence, 275 outcome: DeliveryAttemptOutcome, 276 retry: Option<RetrySchedule>, 277 applied_at_unix_ms: u64, 278 ) -> Result<Self, Error> { 279 if applied_at_unix_ms == 0 { 280 return Err(Error::AtomicWorkflowMismatch); 281 } 282 Ok(Self { 283 plan_id, 284 fence, 285 outcome, 286 retry, 287 applied_at_unix_ms, 288 }) 289 } 290 pub const fn plan_id(&self) -> AuthoredDeliveryPlanId { 291 self.plan_id 292 } 293 pub const fn fence(&self) -> &WorkFence { 294 &self.fence 295 } 296 pub const fn outcome(&self) -> &DeliveryAttemptOutcome { 297 &self.outcome 298 } 299 pub const fn retry(&self) -> Option<&RetrySchedule> { 300 self.retry.as_ref() 301 } 302 pub const fn applied_at_unix_ms(&self) -> u64 { 303 self.applied_at_unix_ms 304 } 305 } 306 307 #[derive(Clone, Debug, Eq, PartialEq)] 308 pub enum AuthoredWorkTarget { 309 Artifact(AuthoredArtifactId), 310 DeliveryPlan(AuthoredDeliveryPlanId), 311 } 312 313 #[derive(Clone, Debug, Eq, PartialEq)] 314 pub enum ClaimAuthoredTarget { 315 ArtifactSigning(AuthoredArtifactId), 316 ArtifactAdmission(AuthoredArtifactId), 317 DeliveryPlan(AuthoredDeliveryPlanId), 318 } 319 320 #[derive(Clone, Debug, Eq, PartialEq)] 321 pub struct ClaimAuthoredWork { 322 target: ClaimAuthoredTarget, 323 claim: crate::authored::WorkClaim, 324 } 325 326 impl ClaimAuthoredWork { 327 pub const fn new(target: ClaimAuthoredTarget, claim: crate::authored::WorkClaim) -> Self { 328 Self { target, claim } 329 } 330 pub const fn target(&self) -> &ClaimAuthoredTarget { 331 &self.target 332 } 333 pub const fn claim(&self) -> &crate::authored::WorkClaim { 334 &self.claim 335 } 336 } 337 338 #[derive(Clone, Debug, Eq, PartialEq)] 339 pub struct ApplyWorkFailure { 340 target: AuthoredWorkTarget, 341 fence: WorkFence, 342 failure: WorkFailure, 343 retry: Option<RetrySchedule>, 344 applied_at_unix_ms: u64, 345 } 346 347 impl ApplyWorkFailure { 348 pub fn new( 349 target: AuthoredWorkTarget, 350 fence: WorkFence, 351 failure: WorkFailure, 352 retry: Option<RetrySchedule>, 353 applied_at_unix_ms: u64, 354 ) -> Result<Self, Error> { 355 if applied_at_unix_ms == 0 { 356 return Err(Error::AtomicWorkflowMismatch); 357 } 358 Ok(Self { 359 target, 360 fence, 361 failure, 362 retry, 363 applied_at_unix_ms, 364 }) 365 } 366 pub const fn target(&self) -> &AuthoredWorkTarget { 367 &self.target 368 } 369 pub const fn fence(&self) -> &WorkFence { 370 &self.fence 371 } 372 pub const fn failure(&self) -> &WorkFailure { 373 &self.failure 374 } 375 pub const fn retry(&self) -> Option<&RetrySchedule> { 376 self.retry.as_ref() 377 } 378 pub const fn applied_at_unix_ms(&self) -> u64 { 379 self.applied_at_unix_ms 380 } 381 } 382 383 #[derive(Clone, Debug, Eq, PartialEq)] 384 pub enum CancelAuthoredTarget { 385 ArtifactSigning(AuthoredArtifactId), 386 ArtifactAdmission(AuthoredArtifactId), 387 DeliveryPlan(AuthoredDeliveryPlanId), 388 } 389 390 #[derive(Clone, Debug, Eq, PartialEq)] 391 pub struct CancelAuthoredWork { 392 target: CancelAuthoredTarget, 393 expected_revision: NonZeroU64, 394 cancelled_at_unix_ms: u64, 395 } 396 397 impl CancelAuthoredWork { 398 pub const fn new( 399 target: CancelAuthoredTarget, 400 expected_revision: NonZeroU64, 401 cancelled_at_unix_ms: u64, 402 ) -> Result<Self, Error> { 403 if cancelled_at_unix_ms == 0 { 404 return Err(Error::AtomicWorkflowMismatch); 405 } 406 Ok(Self { 407 target, 408 expected_revision, 409 cancelled_at_unix_ms, 410 }) 411 } 412 pub const fn target(&self) -> &CancelAuthoredTarget { 413 &self.target 414 } 415 pub const fn expected_revision(&self) -> NonZeroU64 { 416 self.expected_revision 417 } 418 pub const fn cancelled_at_unix_ms(&self) -> u64 { 419 self.cancelled_at_unix_ms 420 } 421 } 422 423 #[derive(Clone, Debug, Eq, PartialEq)] 424 pub enum AuthoredAtomicCommand { 425 Prepare(PrepareAuthoredOperation), 426 PrepareFromDraft(Box<PrepareFromDraft>), 427 Claim(ClaimAuthoredWork), 428 ApplySigned(ApplySignedArtifact), 429 RecordSigned(RecordSignedArtifact), 430 ApplyAdmission(ApplyAdmissionResult), 431 ApplyDelivery(ApplyDeliveryAttempt), 432 RecordDelivery(RecordDeliveryFact), 433 ReconcileDelivery(ReconcileDeliveryFacts), 434 ApplyFailure(ApplyWorkFailure), 435 Cancel(CancelAuthoredWork), 436 } 437 438 impl AuthoredAtomicCommand { 439 pub fn commit_id(&self) -> AtomicCommitId { 440 if let Self::PrepareFromDraft(value) = self { 441 return value.commit_id(); 442 } 443 let digest = self.digest(); 444 let mut hasher = Sha256::new(); 445 hash_field(&mut hasher, b"radroots.authored.atomic.id.v2"); 446 hash_field(&mut hasher, self.phase_bytes()); 447 hash_field(&mut hasher, &self.target_bytes()); 448 if let Some(generation) = self.generation() { 449 hasher.update(generation.get().to_be_bytes()); 450 } 451 hash_field(&mut hasher, digest.as_bytes()); 452 let bytes: [u8; 32] = hasher.finalize().into(); 453 let mut id = [0_u8; 16]; 454 id.copy_from_slice(&bytes[..16]); 455 AtomicCommitId::new(id).expect("SHA-256 derived commit identity is nonzero") 456 } 457 458 pub fn digest(&self) -> AtomicCommitDigest { 459 let mut hasher = Sha256::new(); 460 hash_field(&mut hasher, b"radroots.authored.atomic.digest.v2"); 461 hash_field(&mut hasher, self.phase_bytes()); 462 hash_field(&mut hasher, &self.target_bytes()); 463 match self { 464 Self::PrepareFromDraft(value) => { 465 hash_field(&mut hasher, value.commit_id().as_bytes()); 466 hash_field(&mut hasher, value.source().payload_sha256()); 467 hash_field(&mut hasher, value.intent().payload_sha256()); 468 hash_field(&mut hasher, value.preparation().input_digest().as_bytes()); 469 } 470 Self::Prepare(value) => hash_field(&mut hasher, value.input_digest.as_bytes()), 471 Self::Claim(value) => { 472 hash_field(&mut hasher, value.claim.token()); 473 hasher.update(value.claim.generation().get().to_be_bytes()); 474 hasher.update(value.claim.row_revision().get().to_be_bytes()); 475 } 476 Self::ApplySigned(value) => hash_field(&mut hasher, value.event.raw_json().as_bytes()), 477 Self::RecordSigned(value) => { 478 hash_field(&mut hasher, value.operation_id().as_bytes()); 479 let claim = value.claim(); 480 hash_field(&mut hasher, claim.token()); 481 hash_field(&mut hasher, claim.owner().as_bytes()); 482 hasher.update(claim.generation().get().to_be_bytes()); 483 hasher.update(claim.row_revision().get().to_be_bytes()); 484 hasher.update(claim.acquired_at_unix_ms().to_be_bytes()); 485 hasher.update(claim.expires_at_unix_ms().to_be_bytes()); 486 hash_field(&mut hasher, value.event().raw_json().as_bytes()); 487 } 488 Self::ApplyAdmission(value) => { 489 hasher.update([value.state as u8]); 490 hash_failure(&mut hasher, value.failure.as_ref()); 491 } 492 Self::ApplyDelivery(value) => hash_delivery(&mut hasher, &value.outcome), 493 Self::RecordDelivery(value) => { 494 hash_field(&mut hasher, value.artifact_id().as_bytes()); 495 let claim = value.claim(); 496 hash_field(&mut hasher, claim.token()); 497 hash_field(&mut hasher, claim.owner().as_bytes()); 498 hasher.update(claim.generation().get().to_be_bytes()); 499 hasher.update(claim.row_revision().get().to_be_bytes()); 500 hasher.update(claim.acquired_at_unix_ms().to_be_bytes()); 501 hasher.update(claim.expires_at_unix_ms().to_be_bytes()); 502 value.hash_outcome(&mut hasher); 503 } 504 Self::ReconcileDelivery(value) => value.hash_into(&mut hasher), 505 Self::ApplyFailure(value) => hash_failure(&mut hasher, Some(&value.failure)), 506 Self::Cancel(value) => hasher.update(value.cancelled_at_unix_ms.to_be_bytes()), 507 } 508 AtomicCommitDigest::new(hasher.finalize().into()) 509 } 510 511 pub const fn requested_at_unix_ms(&self) -> u64 { 512 match self { 513 Self::PrepareFromDraft(value) => value.preparation().requested_at_unix_ms(), 514 Self::Prepare(value) => value.requested_at_unix_ms, 515 Self::Claim(value) => value.claim.acquired_at_unix_ms(), 516 Self::ApplySigned(value) => value.applied_at_unix_ms, 517 Self::RecordSigned(value) => value.observed_at_unix_ms(), 518 Self::ApplyAdmission(value) => value.applied_at_unix_ms, 519 Self::ApplyDelivery(value) => value.applied_at_unix_ms, 520 Self::RecordDelivery(value) => value.observed_at_unix_ms(), 521 Self::ReconcileDelivery(value) => value.reconciled_at_unix_ms(), 522 Self::ApplyFailure(value) => value.applied_at_unix_ms, 523 Self::Cancel(value) => value.cancelled_at_unix_ms, 524 } 525 } 526 527 fn phase_bytes(&self) -> &'static [u8] { 528 match self { 529 Self::PrepareFromDraft(_) => b"draft_submission_v1", 530 Self::Prepare(_) => b"prepare", 531 Self::Claim(_) => b"claim", 532 Self::ApplySigned(_) => b"signing", 533 Self::RecordSigned(_) => b"signed_fact_v1", 534 Self::ApplyAdmission(_) => b"admission", 535 Self::ApplyDelivery(_) => b"delivery", 536 Self::RecordDelivery(_) => b"delivery_fact_v1", 537 Self::ReconcileDelivery(_) => b"delivery_reconciliation_v1", 538 Self::ApplyFailure(value) => match value.failure.phase() { 539 WorkPhase::Signing => b"signing_failure", 540 WorkPhase::Admission => b"admission_failure", 541 WorkPhase::Delivery => b"delivery_failure", 542 }, 543 Self::Cancel(_) => b"cancel", 544 } 545 } 546 547 fn target_bytes(&self) -> [u8; 16] { 548 match self { 549 Self::PrepareFromDraft(value) => { 550 *value.preparation().operation().operation_id().as_bytes() 551 } 552 Self::Prepare(value) => *value.operation.operation_id().as_bytes(), 553 Self::Claim(value) => match &value.target { 554 ClaimAuthoredTarget::ArtifactSigning(id) 555 | ClaimAuthoredTarget::ArtifactAdmission(id) => *id.as_bytes(), 556 ClaimAuthoredTarget::DeliveryPlan(id) => *id.as_bytes(), 557 }, 558 Self::ApplySigned(value) => *value.artifact_id.as_bytes(), 559 Self::RecordSigned(value) => *value.artifact_id().as_bytes(), 560 Self::ApplyAdmission(value) => *value.artifact_id.as_bytes(), 561 Self::ApplyDelivery(value) => *value.plan_id.as_bytes(), 562 Self::RecordDelivery(value) => *value.plan_id().as_bytes(), 563 Self::ReconcileDelivery(value) => *value.plan_id().as_bytes(), 564 Self::ApplyFailure(value) => match &value.target { 565 AuthoredWorkTarget::Artifact(id) => *id.as_bytes(), 566 AuthoredWorkTarget::DeliveryPlan(id) => *id.as_bytes(), 567 }, 568 Self::Cancel(value) => match &value.target { 569 CancelAuthoredTarget::ArtifactSigning(id) 570 | CancelAuthoredTarget::ArtifactAdmission(id) => *id.as_bytes(), 571 CancelAuthoredTarget::DeliveryPlan(id) => *id.as_bytes(), 572 }, 573 } 574 } 575 576 fn generation(&self) -> Option<NonZeroU64> { 577 match self { 578 Self::ApplySigned(value) => Some(value.fence.generation), 579 Self::RecordSigned(value) => Some(value.claim().generation()), 580 Self::Claim(value) => Some(value.claim.generation()), 581 Self::ApplyAdmission(value) => Some(value.fence.generation), 582 Self::ApplyDelivery(value) => Some(value.fence.generation), 583 Self::RecordDelivery(value) => Some(value.claim().generation()), 584 Self::ApplyFailure(value) => Some(value.fence.generation), 585 Self::Prepare(_) 586 | Self::PrepareFromDraft(_) 587 | Self::Cancel(_) 588 | Self::ReconcileDelivery(_) => None, 589 } 590 } 591 } 592 593 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] 594 #[cfg_attr(feature = "serde", serde(rename_all = "snake_case"))] 595 #[derive(Clone, Debug, Eq, PartialEq)] 596 pub enum AuthoredAtomicOutcome { 597 Submitted(Box<PrepareFromDraft>), 598 Prepared { 599 operation: AuthoredOperation, 600 artifacts: Vec<AuthoredArtifact>, 601 delivery_plans: Vec<AuthoredDeliveryPlan>, 602 }, 603 Artifact(AuthoredArtifact), 604 DeliveryPlan(AuthoredDeliveryPlan), 605 } 606 607 #[derive(Clone, Debug, Eq, PartialEq)] 608 pub struct AuthoredAtomicReceipt { 609 commit_id: AtomicCommitId, 610 digest: AtomicCommitDigest, 611 disposition: AtomicCommitDisposition, 612 committed_at_unix_ms: u64, 613 outcome: AuthoredAtomicOutcome, 614 } 615 616 impl AuthoredAtomicReceipt { 617 /// Exact submission replay compares the entire validated captured request. 618 pub fn matches_command(&self, command: &AuthoredAtomicCommand) -> bool { 619 if self.commit_id != command.commit_id() || self.digest != command.digest() { 620 return false; 621 } 622 match (command, &self.outcome) { 623 ( 624 AuthoredAtomicCommand::ReconcileDelivery(value), 625 AuthoredAtomicOutcome::DeliveryPlan(plan), 626 ) => value.matches_plan(plan), 627 (AuthoredAtomicCommand::ReconcileDelivery(_), _) => false, 628 ( 629 AuthoredAtomicCommand::RecordDelivery(value), 630 AuthoredAtomicOutcome::DeliveryPlan(plan), 631 ) => value.matches_plan(plan), 632 (AuthoredAtomicCommand::RecordDelivery(_), _) => false, 633 ( 634 AuthoredAtomicCommand::RecordSigned(value), 635 AuthoredAtomicOutcome::Artifact(artifact), 636 ) => signed_fact_matches(value, artifact), 637 (AuthoredAtomicCommand::RecordSigned(_), _) => false, 638 ( 639 AuthoredAtomicCommand::PrepareFromDraft(request), 640 AuthoredAtomicOutcome::Submitted(committed), 641 ) => request == committed, 642 (AuthoredAtomicCommand::PrepareFromDraft(_), _) 643 | (_, AuthoredAtomicOutcome::Submitted(_)) => false, 644 _ => true, 645 } 646 } 647 648 pub fn new( 649 command: &AuthoredAtomicCommand, 650 disposition: AtomicCommitDisposition, 651 committed_at_unix_ms: u64, 652 outcome: AuthoredAtomicOutcome, 653 ) -> Result<Self, Error> { 654 if committed_at_unix_ms < command.requested_at_unix_ms() { 655 return Err(Error::AtomicWorkflowMismatch); 656 } 657 if let (AuthoredAtomicCommand::RecordSigned(_), AuthoredAtomicOutcome::Artifact(artifact)) = 658 (command, &outcome) 659 && committed_at_unix_ms < artifact.updated_at_unix_ms() 660 { 661 return Err(Error::AtomicWorkflowMismatch); 662 } 663 match (command, &outcome) { 664 ( 665 AuthoredAtomicCommand::ReconcileDelivery(value), 666 AuthoredAtomicOutcome::DeliveryPlan(plan), 667 ) if value.matches_plan(plan) && committed_at_unix_ms >= plan.updated_at_unix_ms() => { 668 plan.validate()? 669 } 670 (AuthoredAtomicCommand::ReconcileDelivery(_), _) => { 671 return Err(Error::AtomicWorkflowMismatch); 672 } 673 ( 674 AuthoredAtomicCommand::RecordDelivery(value), 675 AuthoredAtomicOutcome::DeliveryPlan(plan), 676 ) if value.matches_plan(plan) && committed_at_unix_ms >= plan.updated_at_unix_ms() => { 677 plan.validate()? 678 } 679 (AuthoredAtomicCommand::RecordDelivery(_), _) => { 680 return Err(Error::AtomicWorkflowMismatch); 681 } 682 ( 683 AuthoredAtomicCommand::RecordSigned(value), 684 AuthoredAtomicOutcome::Artifact(artifact), 685 ) if signed_fact_matches(value, artifact) => artifact.validate()?, 686 (AuthoredAtomicCommand::RecordSigned(_), _) => { 687 return Err(Error::AtomicWorkflowMismatch); 688 } 689 ( 690 AuthoredAtomicCommand::PrepareFromDraft(request), 691 AuthoredAtomicOutcome::Submitted(value), 692 ) if request == value => value.validate()?, 693 (AuthoredAtomicCommand::PrepareFromDraft(_), _) 694 | (_, AuthoredAtomicOutcome::Submitted(_)) => { 695 return Err(Error::AtomicWorkflowMismatch); 696 } 697 _ => {} 698 } 699 Ok(Self { 700 commit_id: command.commit_id(), 701 digest: command.digest(), 702 disposition, 703 committed_at_unix_ms, 704 outcome, 705 }) 706 } 707 708 pub fn from_durable_parts( 709 commit_id: AtomicCommitId, 710 digest: AtomicCommitDigest, 711 disposition: AtomicCommitDisposition, 712 committed_at_unix_ms: u64, 713 outcome: AuthoredAtomicOutcome, 714 ) -> Result<Self, Error> { 715 if committed_at_unix_ms == 0 || !outcome.is_valid() { 716 return Err(Error::AtomicWorkflowMismatch); 717 } 718 if let AuthoredAtomicOutcome::Submitted(value) = &outcome { 719 let command = AuthoredAtomicCommand::PrepareFromDraft(value.clone()); 720 if commit_id != command.commit_id() 721 || digest != command.digest() 722 || committed_at_unix_ms < command.requested_at_unix_ms() 723 { 724 return Err(Error::AtomicWorkflowMismatch); 725 } 726 } 727 Ok(Self { 728 commit_id, 729 digest, 730 disposition, 731 committed_at_unix_ms, 732 outcome, 733 }) 734 } 735 pub const fn commit_id(&self) -> AtomicCommitId { 736 self.commit_id 737 } 738 pub const fn digest(&self) -> AtomicCommitDigest { 739 self.digest 740 } 741 pub const fn disposition(&self) -> AtomicCommitDisposition { 742 self.disposition 743 } 744 pub const fn committed_at_unix_ms(&self) -> u64 { 745 self.committed_at_unix_ms 746 } 747 pub const fn outcome(&self) -> &AuthoredAtomicOutcome { 748 &self.outcome 749 } 750 } 751 752 fn signed_fact_matches(value: &RecordSignedArtifact, artifact: &AuthoredArtifact) -> bool { 753 artifact.operation_id() == value.operation_id() 754 && artifact.artifact_id() == value.artifact_id() 755 && artifact 756 .signed() 757 .is_some_and(|signed| signed.event() == value.event()) 758 } 759 760 impl AuthoredAtomicOutcome { 761 fn is_valid(&self) -> bool { 762 match self { 763 Self::Prepared { 764 operation, 765 artifacts, 766 delivery_plans, 767 } => { 768 operation.artifact_ids().len() == artifacts.len() 769 && artifacts.iter().enumerate().all(|(ordinal, artifact)| { 770 artifact.operation_id() == operation.operation_id() 771 && operation.artifact_ids().get(ordinal) 772 == Some(&artifact.artifact_id()) 773 && artifact.validate().is_ok() 774 }) 775 && delivery_plans.iter().all(|plan| { 776 artifacts 777 .iter() 778 .any(|artifact| artifact.artifact_id() == plan.artifact_id()) 779 && plan.validate().is_ok() 780 }) 781 } 782 Self::Submitted(value) => value.validate().is_ok(), 783 Self::Artifact(artifact) => artifact.validate().is_ok(), 784 Self::DeliveryPlan(plan) => plan.validate().is_ok(), 785 } 786 } 787 } 788 789 pub trait AuthoredAtomicStorage: Send + Sync { 790 /// Read exact issued-claim provenance; unsupported backends fail closed. 791 fn authored_delivery_history( 792 &self, 793 _plan_id: AuthoredDeliveryPlanId, 794 ) -> BoxFuture<'_, Result<Option<crate::authored_delivery::AuthoredDeliveryHistory>, Error>> 795 { 796 Box::pin(async { Err(Error::BackendUnavailable) }) 797 } 798 fn execute_authored( 799 &self, 800 command: AuthoredAtomicCommand, 801 ) -> BoxFuture<'_, Result<AuthoredAtomicReceipt, Error>>; 802 fn authored_receipt( 803 &self, 804 commit_id: AtomicCommitId, 805 ) -> BoxFuture<'_, Result<Option<AuthoredAtomicReceipt>, Error>>; 806 fn authored_operation( 807 &self, 808 operation_id: OperationInstanceId, 809 ) -> BoxFuture<'_, Result<Option<AuthoredOperation>, Error>>; 810 fn authored_artifact( 811 &self, 812 artifact_id: AuthoredArtifactId, 813 ) -> BoxFuture<'_, Result<Option<AuthoredArtifact>, Error>>; 814 fn authored_delivery_plan( 815 &self, 816 plan_id: AuthoredDeliveryPlanId, 817 ) -> BoxFuture<'_, Result<Option<AuthoredDeliveryPlan>, Error>>; 818 } 819 820 fn hash_failure(hasher: &mut Sha256, failure: Option<&WorkFailure>) { 821 if let Some(failure) = failure { 822 hash_field(hasher, failure.code().as_bytes()); 823 hasher.update([failure.phase() as u8, failure.class() as u8]); 824 hasher.update( 825 failure 826 .retry_after_unix_ms() 827 .unwrap_or_default() 828 .to_be_bytes(), 829 ); 830 if let Some(diagnostic) = failure.diagnostic() { 831 hash_field(hasher, diagnostic.as_bytes()); 832 } 833 } 834 } 835 836 fn hash_delivery(hasher: &mut Sha256, outcome: &DeliveryAttemptOutcome) { 837 let entries = match outcome { 838 DeliveryAttemptOutcome::Receipt(receipt) => receipt.target_receipts(), 839 DeliveryAttemptOutcome::SinkFailure(failure) => { 840 hash_field(hasher, failure.code().as_bytes()); 841 hasher.update([failure.retryability() as u8]); 842 failure.partial_evidence() 843 } 844 }; 845 for entry in entries { 846 hash_field(hasher, entry.target().fingerprint().as_str().as_bytes()); 847 hasher.update([entry.was_attempted() as u8, entry.outcome().kind() as u8]); 848 hasher.update([entry.outcome().retryability() as u8]); 849 if let Some(code) = entry.outcome().code() { 850 hash_field(hasher, code.as_bytes()); 851 } 852 if let Some(message) = entry.outcome().message() { 853 hash_field(hasher, message.as_bytes()); 854 } 855 } 856 } 857 858 fn hash_field(hasher: &mut Sha256, value: &[u8]) { 859 hasher.update(u64::try_from(value.len()).unwrap_or(u64::MAX).to_be_bytes()); 860 hasher.update(value); 861 } 862 863 const fn bytes_are_zero(bytes: &[u8; 16]) -> bool { 864 let mut index = 0; 865 while index < bytes.len() { 866 if bytes[index] != 0 { 867 return false; 868 } 869 index += 1; 870 } 871 true 872 }