journal.rs (20512B)
1 //! Durable operation journal contracts. 2 //! 3 //! [`JournalState::Committed`] is the local durable commit point. Cancellation 4 //! before that point records recoverable work and may be resumed with the same 5 //! idempotency key. Cancellation observed after it never claims rollback: 6 //! callers must receive committed state and may continue pending delivery. 7 8 use core::fmt; 9 pub use radroots_event::EventId; 10 pub use radroots_protocol::runtime::v1::OperationId; 11 pub use radroots_transport::BoxFuture; 12 13 use crate::Error; 14 15 /// Maximum UTF-8 bytes in a journal idempotency key. 16 pub const IDEMPOTENCY_KEY_MAX_BYTES: usize = 256; 17 /// Maximum records returned by one recoverable-work query. 18 pub const RECOVERABLE_QUERY_LIMIT_MAX: u16 = 256; 19 20 /// Host-generated identity for one durable operation execution. 21 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] 22 #[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)] 23 pub struct OperationInstanceId([u8; 16]); 24 25 impl OperationInstanceId { 26 pub const fn new(bytes: [u8; 16]) -> Result<Self, Error> { 27 if all_zero(&bytes) { 28 return Err(Error::InvalidOperationInstanceId); 29 } 30 Ok(Self(bytes)) 31 } 32 33 pub const fn as_bytes(&self) -> &[u8; 16] { 34 &self.0 35 } 36 } 37 38 const fn all_zero(bytes: &[u8; 16]) -> bool { 39 let mut index = 0; 40 while index < bytes.len() { 41 if bytes[index] != 0 { 42 return false; 43 } 44 index += 1; 45 } 46 true 47 } 48 49 /// Validated caller-owned idempotency key. 50 #[derive(Clone, Eq, Hash, Ord, PartialEq, PartialOrd)] 51 pub struct IdempotencyKey(String); 52 53 impl IdempotencyKey { 54 pub fn parse(value: impl AsRef<str>) -> Result<Self, Error> { 55 let value = value.as_ref(); 56 if value.is_empty() 57 || value.len() > IDEMPOTENCY_KEY_MAX_BYTES 58 || value != value.trim() 59 || value.chars().any(char::is_control) 60 { 61 return Err(Error::InvalidIdempotencyKey); 62 } 63 Ok(Self(value.to_owned())) 64 } 65 66 pub fn as_str(&self) -> &str { 67 self.0.as_str() 68 } 69 } 70 71 impl fmt::Debug for IdempotencyKey { 72 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 73 formatter 74 .debug_struct("IdempotencyKey") 75 .field("value", &"[REDACTED]") 76 .field("bytes", &self.0.len()) 77 .finish() 78 } 79 } 80 81 #[cfg(feature = "serde")] 82 impl serde::Serialize for IdempotencyKey { 83 fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error> 84 where 85 S: serde::Serializer, 86 { 87 serializer.serialize_str(self.as_str()) 88 } 89 } 90 91 #[cfg(feature = "serde")] 92 impl<'de> serde::Deserialize<'de> for IdempotencyKey { 93 fn deserialize<D>(deserializer: D) -> Result<Self, D::Error> 94 where 95 D: serde::Deserializer<'de>, 96 { 97 let value = <String as serde::Deserialize>::deserialize(deserializer)?; 98 Self::parse(value).map_err(serde::de::Error::custom) 99 } 100 } 101 102 /// SHA-256 digest of canonical operation input, computed by its domain owner. 103 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] 104 #[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)] 105 pub struct IdempotencyDigest([u8; 32]); 106 107 impl IdempotencyDigest { 108 pub const fn new(bytes: [u8; 32]) -> Self { 109 Self(bytes) 110 } 111 112 pub const fn as_bytes(&self) -> &[u8; 32] { 113 &self.0 114 } 115 } 116 117 /// Non-zero optimistic journal revision. 118 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] 119 #[derive(Clone, Copy, Debug, Eq, Ord, PartialEq, PartialOrd)] 120 pub struct JournalRevision(u64); 121 122 impl JournalRevision { 123 pub const INITIAL: Self = Self(1); 124 125 pub const fn new(value: u64) -> Result<Self, Error> { 126 if value == 0 { 127 return Err(Error::InvalidJournalRevision); 128 } 129 Ok(Self(value)) 130 } 131 132 pub const fn get(self) -> u64 { 133 self.0 134 } 135 136 fn next(self) -> Result<Self, Error> { 137 self.0 138 .checked_add(1) 139 .map(Self) 140 .ok_or(Error::CorruptJournalRecord) 141 } 142 } 143 144 /// Durable lifecycle stage. 145 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] 146 #[cfg_attr(feature = "serde", serde(rename_all = "snake_case"))] 147 #[derive(Clone, Copy, Debug, Eq, PartialEq)] 148 pub enum JournalStage { 149 Prepared, 150 Signed, 151 Recoverable, 152 Committed, 153 } 154 155 /// Point from which recoverable work resumes. 156 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] 157 #[cfg_attr(feature = "serde", serde(rename_all = "snake_case"))] 158 #[derive(Clone, Debug, Eq, PartialEq)] 159 pub enum RecoveryPoint { 160 Prepared, 161 Signed { event_id: EventId }, 162 } 163 164 /// Stable class of recoverable interruption. 165 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] 166 #[cfg_attr(feature = "serde", serde(rename_all = "snake_case"))] 167 #[derive(Clone, Copy, Debug, Eq, PartialEq)] 168 pub enum RecoveryReason { 169 CancelledBeforeCommit, 170 SignerUnavailable, 171 TransportUnavailable, 172 StorageUnavailable, 173 DeadlineExceeded, 174 Interrupted, 175 } 176 177 /// Durable recovery evidence without backend or secret detail. 178 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] 179 #[derive(Clone, Debug, Eq, PartialEq)] 180 pub struct RecoveryRecord { 181 point: RecoveryPoint, 182 reason: RecoveryReason, 183 attempt: u32, 184 retry_not_before_unix_ms: Option<u64>, 185 } 186 187 impl RecoveryRecord { 188 pub const fn new( 189 point: RecoveryPoint, 190 reason: RecoveryReason, 191 attempt: u32, 192 retry_not_before_unix_ms: Option<u64>, 193 ) -> Result<Self, Error> { 194 if attempt == 0 { 195 return Err(Error::InvalidRecoveryAttempt); 196 } 197 if matches!(retry_not_before_unix_ms, Some(0)) { 198 return Err(Error::InvalidRecoveryDeadline); 199 } 200 Ok(Self { 201 point, 202 reason, 203 attempt, 204 retry_not_before_unix_ms, 205 }) 206 } 207 208 pub const fn point(&self) -> &RecoveryPoint { 209 &self.point 210 } 211 212 pub const fn reason(&self) -> RecoveryReason { 213 self.reason 214 } 215 216 pub const fn attempt(&self) -> u32 { 217 self.attempt 218 } 219 220 pub const fn retry_not_before_unix_ms(&self) -> Option<u64> { 221 self.retry_not_before_unix_ms 222 } 223 } 224 225 /// Cancellation observation relative to the local durable commit point. 226 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] 227 #[cfg_attr(feature = "serde", serde(rename_all = "snake_case"))] 228 #[derive(Clone, Copy, Debug, Eq, PartialEq)] 229 pub enum CancellationState { 230 NotRequested, 231 CancelledBeforeCommit, 232 ObservedAfterCommit, 233 } 234 235 /// State-specific durable journal data. 236 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] 237 #[cfg_attr(feature = "serde", serde(rename_all = "snake_case"))] 238 #[derive(Clone, Debug, Eq, PartialEq)] 239 pub enum JournalState { 240 Prepared, 241 Signed { 242 event_id: EventId, 243 }, 244 Recoverable(RecoveryRecord), 245 /// Local state is durably committed; remote delivery may remain pending. 246 Committed { 247 event_id: EventId, 248 committed_at_unix_ms: u64, 249 }, 250 } 251 252 impl JournalState { 253 pub const fn stage(&self) -> JournalStage { 254 match self { 255 Self::Prepared => JournalStage::Prepared, 256 Self::Signed { .. } => JournalStage::Signed, 257 Self::Recoverable(_) => JournalStage::Recoverable, 258 Self::Committed { .. } => JournalStage::Committed, 259 } 260 } 261 } 262 263 /// Durable operation record. 264 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] 265 #[derive(Clone, Debug, Eq, PartialEq)] 266 pub struct OperationRecord { 267 instance_id: OperationInstanceId, 268 operation_id: OperationId, 269 idempotency_key: IdempotencyKey, 270 input_digest: IdempotencyDigest, 271 prepared_at_unix_ms: u64, 272 revision: JournalRevision, 273 state: JournalState, 274 cancellation: CancellationState, 275 } 276 277 impl OperationRecord { 278 #[allow(clippy::too_many_arguments)] 279 pub fn from_parts( 280 instance_id: OperationInstanceId, 281 operation_id: OperationId, 282 idempotency_key: IdempotencyKey, 283 input_digest: IdempotencyDigest, 284 prepared_at_unix_ms: u64, 285 revision: JournalRevision, 286 state: JournalState, 287 cancellation: CancellationState, 288 ) -> Result<Self, Error> { 289 if prepared_at_unix_ms == 0 { 290 return Err(Error::InvalidOperationTimestamp); 291 } 292 validate_state(prepared_at_unix_ms, &state, cancellation)?; 293 Ok(Self { 294 instance_id, 295 operation_id, 296 idempotency_key, 297 input_digest, 298 prepared_at_unix_ms, 299 revision, 300 state, 301 cancellation, 302 }) 303 } 304 305 pub const fn instance_id(&self) -> OperationInstanceId { 306 self.instance_id 307 } 308 309 pub const fn operation_id(&self) -> OperationId { 310 self.operation_id 311 } 312 313 pub const fn idempotency_key(&self) -> &IdempotencyKey { 314 &self.idempotency_key 315 } 316 317 pub const fn input_digest(&self) -> IdempotencyDigest { 318 self.input_digest 319 } 320 321 pub const fn prepared_at_unix_ms(&self) -> u64 { 322 self.prepared_at_unix_ms 323 } 324 325 pub const fn revision(&self) -> JournalRevision { 326 self.revision 327 } 328 329 pub const fn state(&self) -> &JournalState { 330 &self.state 331 } 332 333 pub const fn cancellation(&self) -> CancellationState { 334 self.cancellation 335 } 336 337 /// Applies one optimistic transition without allowing lifecycle regressions. 338 pub fn transition(&self, transition: &JournalTransition) -> Result<Self, Error> { 339 if transition.instance_id != self.instance_id { 340 return Err(Error::OperationIdentityMismatch); 341 } 342 if transition.expected_revision != self.revision { 343 return Err(Error::JournalRevisionConflict); 344 } 345 346 let (state, cancellation) = apply_transition(self, &transition.kind)?; 347 Self::from_parts( 348 self.instance_id, 349 self.operation_id, 350 self.idempotency_key.clone(), 351 self.input_digest, 352 self.prepared_at_unix_ms, 353 self.revision.next()?, 354 state, 355 cancellation, 356 ) 357 } 358 } 359 360 fn validate_state( 361 prepared_at: u64, 362 state: &JournalState, 363 cancellation: CancellationState, 364 ) -> Result<(), Error> { 365 if let JournalState::Committed { 366 committed_at_unix_ms, 367 .. 368 } = state 369 && (*committed_at_unix_ms == 0 || *committed_at_unix_ms < prepared_at) 370 { 371 return Err(Error::CorruptJournalRecord); 372 } 373 if let JournalState::Recoverable(recovery) = state 374 && recovery 375 .retry_not_before_unix_ms() 376 .is_some_and(|deadline| deadline < prepared_at) 377 { 378 return Err(Error::CorruptJournalRecord); 379 } 380 match state { 381 JournalState::Committed { .. } 382 if cancellation == CancellationState::CancelledBeforeCommit => 383 { 384 Err(Error::CorruptJournalRecord) 385 } 386 JournalState::Recoverable(_) if cancellation == CancellationState::ObservedAfterCommit => { 387 Err(Error::CorruptJournalRecord) 388 } 389 JournalState::Recoverable(recovery) 390 if (recovery.reason() == RecoveryReason::CancelledBeforeCommit) 391 != (cancellation == CancellationState::CancelledBeforeCommit) => 392 { 393 Err(Error::CorruptJournalRecord) 394 } 395 JournalState::Prepared | JournalState::Signed { .. } 396 if cancellation != CancellationState::NotRequested => 397 { 398 Err(Error::CorruptJournalRecord) 399 } 400 _ => Ok(()), 401 } 402 } 403 404 fn apply_transition( 405 record: &OperationRecord, 406 transition: &JournalTransitionKind, 407 ) -> Result<(JournalState, CancellationState), Error> { 408 match (record.state(), transition) { 409 (JournalState::Prepared, JournalTransitionKind::Signed { event_id }) 410 if record.cancellation() == CancellationState::NotRequested => 411 { 412 Ok(( 413 JournalState::Signed { 414 event_id: *event_id, 415 }, 416 CancellationState::NotRequested, 417 )) 418 } 419 (JournalState::Committed { .. }, JournalTransitionKind::Cancelled { observed_at }) 420 if *observed_at >= record.prepared_at_unix_ms() => 421 { 422 Ok(( 423 record.state().clone(), 424 CancellationState::ObservedAfterCommit, 425 )) 426 } 427 (JournalState::Committed { .. }, _) => Err(Error::JournalOperationCommitted), 428 (_, JournalTransitionKind::Recoverable { record: recovery }) => Ok(( 429 JournalState::Recoverable(recovery.clone()), 430 if recovery.reason() == RecoveryReason::CancelledBeforeCommit { 431 CancellationState::CancelledBeforeCommit 432 } else { 433 CancellationState::NotRequested 434 }, 435 )), 436 (JournalState::Recoverable(recovery), JournalTransitionKind::Resume) => Ok(( 437 match recovery.point() { 438 RecoveryPoint::Prepared => JournalState::Prepared, 439 RecoveryPoint::Signed { event_id } => JournalState::Signed { 440 event_id: *event_id, 441 }, 442 }, 443 CancellationState::NotRequested, 444 )), 445 ( 446 JournalState::Signed { 447 event_id: signed_id, 448 }, 449 JournalTransitionKind::Committed { 450 event_id, 451 committed_at, 452 }, 453 ) if signed_id == event_id && *committed_at >= record.prepared_at_unix_ms() => Ok(( 454 JournalState::Committed { 455 event_id: *event_id, 456 committed_at_unix_ms: *committed_at, 457 }, 458 CancellationState::NotRequested, 459 )), 460 (JournalState::Prepared, JournalTransitionKind::Cancelled { observed_at }) 461 if *observed_at >= record.prepared_at_unix_ms() => 462 { 463 cancelled(RecoveryPoint::Prepared) 464 } 465 (JournalState::Signed { event_id }, JournalTransitionKind::Cancelled { observed_at }) 466 if *observed_at >= record.prepared_at_unix_ms() => 467 { 468 cancelled(RecoveryPoint::Signed { 469 event_id: *event_id, 470 }) 471 } 472 _ => Err(Error::InvalidJournalTransition), 473 } 474 } 475 476 fn cancelled(point: RecoveryPoint) -> Result<(JournalState, CancellationState), Error> { 477 Ok(( 478 JournalState::Recoverable(RecoveryRecord::new( 479 point, 480 RecoveryReason::CancelledBeforeCommit, 481 1, 482 None, 483 )?), 484 CancellationState::CancelledBeforeCommit, 485 )) 486 } 487 488 /// Input for an idempotent prepare operation. 489 #[derive(Clone, Debug, Eq, PartialEq)] 490 pub struct PrepareOperation { 491 instance_id: OperationInstanceId, 492 operation_id: OperationId, 493 idempotency_key: IdempotencyKey, 494 input_digest: IdempotencyDigest, 495 prepared_at_unix_ms: u64, 496 } 497 498 impl PrepareOperation { 499 pub fn new( 500 instance_id: OperationInstanceId, 501 operation_id: OperationId, 502 idempotency_key: IdempotencyKey, 503 input_digest: IdempotencyDigest, 504 prepared_at_unix_ms: u64, 505 ) -> Result<Self, Error> { 506 if prepared_at_unix_ms == 0 { 507 return Err(Error::InvalidOperationTimestamp); 508 } 509 Ok(Self { 510 instance_id, 511 operation_id, 512 idempotency_key, 513 input_digest, 514 prepared_at_unix_ms, 515 }) 516 } 517 518 pub const fn instance_id(&self) -> OperationInstanceId { 519 self.instance_id 520 } 521 522 pub const fn operation_id(&self) -> OperationId { 523 self.operation_id 524 } 525 526 pub const fn idempotency_key(&self) -> &IdempotencyKey { 527 &self.idempotency_key 528 } 529 530 pub const fn input_digest(&self) -> IdempotencyDigest { 531 self.input_digest 532 } 533 534 pub fn into_record(self) -> Result<OperationRecord, Error> { 535 OperationRecord::from_parts( 536 self.instance_id, 537 self.operation_id, 538 self.idempotency_key, 539 self.input_digest, 540 self.prepared_at_unix_ms, 541 JournalRevision::INITIAL, 542 JournalState::Prepared, 543 CancellationState::NotRequested, 544 ) 545 } 546 } 547 548 /// Result of preparing an idempotent operation. 549 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] 550 #[cfg_attr(feature = "serde", serde(rename_all = "snake_case"))] 551 #[derive(Clone, Copy, Debug, Eq, PartialEq)] 552 pub enum PrepareDisposition { 553 Created, 554 Replay, 555 } 556 557 #[derive(Clone, Debug, Eq, PartialEq)] 558 pub struct PrepareReceipt { 559 disposition: PrepareDisposition, 560 record: OperationRecord, 561 } 562 563 impl PrepareReceipt { 564 pub const fn new(disposition: PrepareDisposition, record: OperationRecord) -> Self { 565 Self { 566 disposition, 567 record, 568 } 569 } 570 571 pub const fn disposition(&self) -> PrepareDisposition { 572 self.disposition 573 } 574 575 pub const fn record(&self) -> &OperationRecord { 576 &self.record 577 } 578 } 579 580 /// Validated optimistic transition request. 581 #[derive(Clone, Debug, Eq, PartialEq)] 582 pub struct JournalTransition { 583 instance_id: OperationInstanceId, 584 expected_revision: JournalRevision, 585 kind: JournalTransitionKind, 586 } 587 588 #[derive(Clone, Debug, Eq, PartialEq)] 589 enum JournalTransitionKind { 590 Signed { 591 event_id: EventId, 592 }, 593 Recoverable { 594 record: RecoveryRecord, 595 }, 596 Resume, 597 Committed { 598 event_id: EventId, 599 committed_at: u64, 600 }, 601 Cancelled { 602 observed_at: u64, 603 }, 604 } 605 606 impl JournalTransition { 607 pub const fn signed( 608 instance_id: OperationInstanceId, 609 expected_revision: JournalRevision, 610 event_id: EventId, 611 ) -> Self { 612 Self { 613 instance_id, 614 expected_revision, 615 kind: JournalTransitionKind::Signed { event_id }, 616 } 617 } 618 619 pub const fn recoverable( 620 instance_id: OperationInstanceId, 621 expected_revision: JournalRevision, 622 record: RecoveryRecord, 623 ) -> Self { 624 Self { 625 instance_id, 626 expected_revision, 627 kind: JournalTransitionKind::Recoverable { record }, 628 } 629 } 630 631 pub const fn resume( 632 instance_id: OperationInstanceId, 633 expected_revision: JournalRevision, 634 ) -> Self { 635 Self { 636 instance_id, 637 expected_revision, 638 kind: JournalTransitionKind::Resume, 639 } 640 } 641 642 pub const fn committed( 643 instance_id: OperationInstanceId, 644 expected_revision: JournalRevision, 645 event_id: EventId, 646 committed_at: u64, 647 ) -> Self { 648 Self { 649 instance_id, 650 expected_revision, 651 kind: JournalTransitionKind::Committed { 652 event_id, 653 committed_at, 654 }, 655 } 656 } 657 658 pub const fn cancelled( 659 instance_id: OperationInstanceId, 660 expected_revision: JournalRevision, 661 observed_at: u64, 662 ) -> Self { 663 Self { 664 instance_id, 665 expected_revision, 666 kind: JournalTransitionKind::Cancelled { observed_at }, 667 } 668 } 669 670 pub const fn instance_id(&self) -> OperationInstanceId { 671 self.instance_id 672 } 673 } 674 675 /// Backend-neutral durable operation journal SPI. 676 pub trait Journal: Send + Sync { 677 /// Creates a record or replays the exact existing operation. A reused key 678 /// with a different operation kind, instance, or digest is a conflict. 679 fn prepare(&self, operation: PrepareOperation) -> BoxFuture<'_, Result<PrepareReceipt, Error>>; 680 681 fn operation( 682 &self, 683 instance_id: OperationInstanceId, 684 ) -> BoxFuture<'_, Result<Option<OperationRecord>, Error>>; 685 686 fn by_idempotency_key( 687 &self, 688 operation_id: OperationId, 689 idempotency_key: IdempotencyKey, 690 ) -> BoxFuture<'_, Result<Option<OperationRecord>, Error>>; 691 692 /// Applies one lifecycle transition atomically at its expected revision. 693 fn transition( 694 &self, 695 transition: JournalTransition, 696 ) -> BoxFuture<'_, Result<OperationRecord, Error>>; 697 698 /// Returns recoverable records in backend-stable order. 699 fn recoverable(&self, limit: u16) -> BoxFuture<'_, Result<Vec<OperationRecord>, Error>>; 700 }