journal.rs (19017B)
1 use futures_executor::block_on; 2 use radroots_event::EventId; 3 use radroots_protocol::runtime::v1::OperationId; 4 use radroots_storage::{ 5 Error, 6 journal::{ 7 CancellationState, IdempotencyDigest, IdempotencyKey, Journal, JournalRevision, 8 JournalStage, JournalState, JournalTransition, OperationInstanceId, OperationRecord, 9 PrepareDisposition, PrepareOperation, PrepareReceipt, RECOVERABLE_QUERY_LIMIT_MAX, 10 RecoveryPoint, RecoveryReason, RecoveryRecord, 11 }, 12 }; 13 use radroots_transport::BoxFuture; 14 use std::sync::Mutex; 15 16 struct MemoryJournal(Mutex<Vec<OperationRecord>>); 17 18 impl MemoryJournal { 19 fn new() -> Self { 20 Self(Mutex::new(Vec::new())) 21 } 22 } 23 24 impl Journal for MemoryJournal { 25 fn prepare(&self, operation: PrepareOperation) -> BoxFuture<'_, Result<PrepareReceipt, Error>> { 26 Box::pin(async move { 27 let mut records = self.0.lock().expect("test journal lock"); 28 if let Some(record) = records 29 .iter() 30 .find(|record| record.idempotency_key() == operation.idempotency_key()) 31 { 32 if record.operation_id() != operation.operation_id() 33 || record.input_digest() != operation.input_digest() 34 || record.instance_id() != operation.instance_id() 35 { 36 return Err(Error::IdempotencyConflict); 37 } 38 return Ok(PrepareReceipt::new( 39 PrepareDisposition::Replay, 40 record.clone(), 41 )); 42 } 43 if records 44 .iter() 45 .any(|record| record.instance_id() == operation.instance_id()) 46 { 47 return Err(Error::OperationIdentityMismatch); 48 } 49 let record = operation.into_record()?; 50 records.push(record.clone()); 51 Ok(PrepareReceipt::new(PrepareDisposition::Created, record)) 52 }) 53 } 54 55 fn operation( 56 &self, 57 instance_id: OperationInstanceId, 58 ) -> BoxFuture<'_, Result<Option<OperationRecord>, Error>> { 59 Box::pin(async move { 60 Ok(self 61 .0 62 .lock() 63 .expect("test journal lock") 64 .iter() 65 .find(|record| record.instance_id() == instance_id) 66 .cloned()) 67 }) 68 } 69 70 fn by_idempotency_key( 71 &self, 72 operation_id: OperationId, 73 idempotency_key: IdempotencyKey, 74 ) -> BoxFuture<'_, Result<Option<OperationRecord>, Error>> { 75 Box::pin(async move { 76 Ok(self 77 .0 78 .lock() 79 .expect("test journal lock") 80 .iter() 81 .find(|record| { 82 record.operation_id() == operation_id 83 && record.idempotency_key() == &idempotency_key 84 }) 85 .cloned()) 86 }) 87 } 88 89 fn transition( 90 &self, 91 transition: JournalTransition, 92 ) -> BoxFuture<'_, Result<OperationRecord, Error>> { 93 Box::pin(async move { 94 let mut records = self.0.lock().expect("test journal lock"); 95 let record = records 96 .iter_mut() 97 .find(|record| record.instance_id() == transition.instance_id()) 98 .ok_or(Error::OperationNotFound)?; 99 let next = record.transition(&transition)?; 100 *record = next.clone(); 101 Ok(next) 102 }) 103 } 104 105 fn recoverable(&self, limit: u16) -> BoxFuture<'_, Result<Vec<OperationRecord>, Error>> { 106 Box::pin(async move { 107 if limit == 0 || limit > RECOVERABLE_QUERY_LIMIT_MAX { 108 return Err(Error::InvalidJournalQueryLimit); 109 } 110 Ok(self 111 .0 112 .lock() 113 .expect("test journal lock") 114 .iter() 115 .filter(|record| record.state().stage() == JournalStage::Recoverable) 116 .take(usize::from(limit)) 117 .cloned() 118 .collect()) 119 }) 120 } 121 } 122 123 fn instance(byte: u8) -> OperationInstanceId { 124 OperationInstanceId::new([byte; 16]).expect("operation instance") 125 } 126 127 fn key(byte: u8) -> IdempotencyKey { 128 IdempotencyKey::parse(format!("sync-push-{byte:02x}")).expect("idempotency key") 129 } 130 131 fn event_id(byte: &str) -> EventId { 132 EventId::parse(byte.repeat(64)).expect("event id") 133 } 134 135 fn prepare(instance_id: OperationInstanceId, digest: u8, at: u64) -> PrepareOperation { 136 PrepareOperation::new( 137 instance_id, 138 OperationId::SyncPush, 139 key(instance_id.as_bytes()[0]), 140 IdempotencyDigest::new([digest; 32]), 141 at, 142 ) 143 .expect("prepare operation") 144 } 145 146 #[test] 147 fn prepare_replays_exact_input_and_rejects_conflicts() { 148 let journal = MemoryJournal::new(); 149 let dynamic: &dyn Journal = &journal; 150 let operation = prepare(instance(1), 2, 100); 151 let created = block_on(dynamic.prepare(operation.clone())).expect("created"); 152 assert_eq!(created.disposition(), PrepareDisposition::Created); 153 assert_eq!(created.record().revision(), JournalRevision::INITIAL); 154 155 let replay = block_on(journal.prepare(operation)).expect("replay"); 156 assert_eq!(replay.disposition(), PrepareDisposition::Replay); 157 assert_eq!(replay.record(), created.record()); 158 assert_eq!( 159 block_on(journal.prepare(prepare(instance(1), 3, 100))), 160 Err(Error::IdempotencyConflict) 161 ); 162 let conflicting_kind = PrepareOperation::new( 163 instance(1), 164 OperationId::FarmPublish, 165 key(1), 166 IdempotencyDigest::new([2; 32]), 167 100, 168 ) 169 .expect("conflicting operation kind"); 170 assert_eq!( 171 block_on(journal.prepare(conflicting_kind)), 172 Err(Error::IdempotencyConflict) 173 ); 174 175 let debug = format!("{:?}", key(1)); 176 assert!(debug.contains("[REDACTED]")); 177 assert!(!debug.contains("sync-push-01")); 178 } 179 180 #[test] 181 fn lifecycle_is_optimistic_monotonic_and_commit_bound() { 182 let journal = MemoryJournal::new(); 183 let instance_id = instance(4); 184 let event_id = event_id("a"); 185 let prepared = block_on(journal.prepare(prepare(instance_id, 5, 100))) 186 .expect("prepare") 187 .record() 188 .clone(); 189 let signed = block_on(journal.transition(JournalTransition::signed( 190 instance_id, 191 prepared.revision(), 192 event_id, 193 ))) 194 .expect("signed"); 195 assert_eq!(signed.state().stage(), JournalStage::Signed); 196 assert_eq!( 197 block_on(journal.transition(JournalTransition::signed( 198 instance_id, 199 prepared.revision(), 200 event_id, 201 ))), 202 Err(Error::JournalRevisionConflict) 203 ); 204 205 let committed = block_on(journal.transition(JournalTransition::committed( 206 instance_id, 207 signed.revision(), 208 event_id, 209 150, 210 ))) 211 .expect("committed"); 212 assert_eq!(committed.state().stage(), JournalStage::Committed); 213 assert_eq!( 214 block_on( 215 journal.transition(JournalTransition::recoverable( 216 instance_id, 217 committed.revision(), 218 RecoveryRecord::new( 219 RecoveryPoint::Signed { event_id }, 220 RecoveryReason::TransportUnavailable, 221 1, 222 None, 223 ) 224 .expect("recovery"), 225 )) 226 ), 227 Err(Error::JournalOperationCommitted) 228 ); 229 } 230 231 #[test] 232 fn cancellation_before_commit_recovers_and_after_commit_preserves_commit() { 233 let journal = MemoryJournal::new(); 234 let before_id = instance(6); 235 let prepared = block_on(journal.prepare(prepare(before_id, 7, 100))) 236 .expect("prepare before") 237 .record() 238 .clone(); 239 let cancelled = block_on(journal.transition(JournalTransition::cancelled( 240 before_id, 241 prepared.revision(), 242 110, 243 ))) 244 .expect("cancel before commit"); 245 assert_eq!(cancelled.state().stage(), JournalStage::Recoverable); 246 assert_eq!( 247 cancelled.cancellation(), 248 CancellationState::CancelledBeforeCommit 249 ); 250 assert_eq!( 251 block_on(journal.recoverable(10)).expect("recovery").len(), 252 1 253 ); 254 255 let resumed = 256 block_on(journal.transition(JournalTransition::resume(before_id, cancelled.revision()))) 257 .expect("resume"); 258 assert_eq!(resumed.state(), &JournalState::Prepared); 259 assert_eq!(resumed.cancellation(), CancellationState::NotRequested); 260 261 let after_id = instance(8); 262 let event_id = event_id("b"); 263 let prepared = block_on(journal.prepare(prepare(after_id, 9, 200))) 264 .expect("prepare after") 265 .record() 266 .clone(); 267 let signed = block_on(journal.transition(JournalTransition::signed( 268 after_id, 269 prepared.revision(), 270 event_id, 271 ))) 272 .expect("signed after"); 273 let committed = block_on(journal.transition(JournalTransition::committed( 274 after_id, 275 signed.revision(), 276 event_id, 277 220, 278 ))) 279 .expect("committed after"); 280 let observed = block_on(journal.transition(JournalTransition::cancelled( 281 after_id, 282 committed.revision(), 283 230, 284 ))) 285 .expect("cancel observed after commit"); 286 assert_eq!(observed.state().stage(), JournalStage::Committed); 287 assert_eq!( 288 observed.cancellation(), 289 CancellationState::ObservedAfterCommit 290 ); 291 } 292 293 #[test] 294 fn invalid_records_and_inputs_fail_closed() { 295 assert_eq!( 296 OperationInstanceId::new([0; 16]), 297 Err(Error::InvalidOperationInstanceId) 298 ); 299 assert_eq!( 300 IdempotencyKey::parse(" bad-key"), 301 Err(Error::InvalidIdempotencyKey) 302 ); 303 assert_eq!( 304 RecoveryRecord::new( 305 RecoveryPoint::Prepared, 306 RecoveryReason::Interrupted, 307 0, 308 None, 309 ), 310 Err(Error::InvalidRecoveryAttempt) 311 ); 312 } 313 314 #[test] 315 fn journal_value_and_state_validation_matrix_is_complete() { 316 let instance_id = instance(1); 317 assert_eq!(instance_id.as_bytes(), &[1; 16]); 318 for invalid in ["", " leading", "trailing ", "bad\nkey"] { 319 assert_eq!( 320 IdempotencyKey::parse(invalid), 321 Err(Error::InvalidIdempotencyKey) 322 ); 323 } 324 assert_eq!( 325 IdempotencyKey::parse("x".repeat(radroots_storage::journal::IDEMPOTENCY_KEY_MAX_BYTES + 1)), 326 Err(Error::InvalidIdempotencyKey) 327 ); 328 assert_eq!( 329 IdempotencyKey::parse("x".repeat(4 * 1024 * 1024)), 330 Err(Error::InvalidIdempotencyKey) 331 ); 332 let maximum = 333 IdempotencyKey::parse("x".repeat(radroots_storage::journal::IDEMPOTENCY_KEY_MAX_BYTES)) 334 .expect("exact maximum key"); 335 assert_eq!( 336 maximum.as_str().len(), 337 radroots_storage::journal::IDEMPOTENCY_KEY_MAX_BYTES 338 ); 339 let idempotency_key = key(1); 340 assert_eq!(idempotency_key.as_str(), "sync-push-01"); 341 let digest = IdempotencyDigest::new([2; 32]); 342 assert_eq!(digest.as_bytes(), &[2; 32]); 343 assert_eq!(JournalRevision::new(0), Err(Error::InvalidJournalRevision)); 344 assert_eq!(JournalRevision::new(2).unwrap().get(), 2); 345 assert_eq!( 346 RecoveryRecord::new( 347 RecoveryPoint::Prepared, 348 RecoveryReason::Interrupted, 349 1, 350 Some(0), 351 ), 352 Err(Error::InvalidRecoveryDeadline) 353 ); 354 let recovery = RecoveryRecord::new( 355 RecoveryPoint::Signed { 356 event_id: event_id("a"), 357 }, 358 RecoveryReason::TransportUnavailable, 359 2, 360 Some(150), 361 ) 362 .unwrap(); 363 assert!(matches!(recovery.point(), RecoveryPoint::Signed { .. })); 364 assert_eq!(recovery.reason(), RecoveryReason::TransportUnavailable); 365 assert_eq!(recovery.attempt(), 2); 366 assert_eq!(recovery.retry_not_before_unix_ms(), Some(150)); 367 for state in [ 368 JournalState::Prepared, 369 JournalState::Signed { 370 event_id: event_id("a"), 371 }, 372 JournalState::Recoverable(recovery.clone()), 373 JournalState::Committed { 374 event_id: event_id("a"), 375 committed_at_unix_ms: 150, 376 }, 377 ] { 378 let expected = match state { 379 JournalState::Prepared => JournalStage::Prepared, 380 JournalState::Signed { .. } => JournalStage::Signed, 381 JournalState::Recoverable(_) => JournalStage::Recoverable, 382 JournalState::Committed { .. } => JournalStage::Committed, 383 }; 384 assert_eq!(state.stage(), expected); 385 } 386 387 let build = |state, cancellation| { 388 OperationRecord::from_parts( 389 instance_id, 390 OperationId::SyncPush, 391 idempotency_key.clone(), 392 digest, 393 100, 394 JournalRevision::INITIAL, 395 state, 396 cancellation, 397 ) 398 }; 399 assert_eq!( 400 OperationRecord::from_parts( 401 instance_id, 402 OperationId::SyncPush, 403 idempotency_key.clone(), 404 digest, 405 0, 406 JournalRevision::INITIAL, 407 JournalState::Prepared, 408 CancellationState::NotRequested, 409 ), 410 Err(Error::InvalidOperationTimestamp) 411 ); 412 for result in [ 413 build( 414 JournalState::Committed { 415 event_id: event_id("a"), 416 committed_at_unix_ms: 0, 417 }, 418 CancellationState::NotRequested, 419 ), 420 build( 421 JournalState::Committed { 422 event_id: event_id("a"), 423 committed_at_unix_ms: 99, 424 }, 425 CancellationState::NotRequested, 426 ), 427 build( 428 JournalState::Recoverable( 429 RecoveryRecord::new( 430 RecoveryPoint::Prepared, 431 RecoveryReason::Interrupted, 432 1, 433 Some(99), 434 ) 435 .unwrap(), 436 ), 437 CancellationState::NotRequested, 438 ), 439 build( 440 JournalState::Committed { 441 event_id: event_id("a"), 442 committed_at_unix_ms: 100, 443 }, 444 CancellationState::CancelledBeforeCommit, 445 ), 446 build( 447 JournalState::Recoverable(recovery.clone()), 448 CancellationState::ObservedAfterCommit, 449 ), 450 build( 451 JournalState::Recoverable(recovery.clone()), 452 CancellationState::CancelledBeforeCommit, 453 ), 454 build( 455 JournalState::Prepared, 456 CancellationState::ObservedAfterCommit, 457 ), 458 build( 459 JournalState::Signed { 460 event_id: event_id("a"), 461 }, 462 CancellationState::CancelledBeforeCommit, 463 ), 464 ] { 465 assert_eq!(result, Err(Error::CorruptJournalRecord)); 466 } 467 468 let operation = prepare(instance_id, 2, 100); 469 assert_eq!(operation.instance_id(), instance_id); 470 assert_eq!(operation.operation_id(), OperationId::SyncPush); 471 assert_eq!(operation.idempotency_key(), &idempotency_key); 472 assert_eq!(operation.input_digest(), digest); 473 let record = operation.into_record().unwrap(); 474 assert_eq!(record.instance_id(), instance_id); 475 assert_eq!(record.operation_id(), OperationId::SyncPush); 476 assert_eq!(record.idempotency_key(), &idempotency_key); 477 assert_eq!(record.input_digest(), digest); 478 assert_eq!(record.prepared_at_unix_ms(), 100); 479 assert_eq!(record.revision(), JournalRevision::INITIAL); 480 assert_eq!(record.cancellation(), CancellationState::NotRequested); 481 let receipt = PrepareReceipt::new(PrepareDisposition::Created, record.clone()); 482 assert_eq!(receipt.disposition(), PrepareDisposition::Created); 483 assert_eq!(receipt.record(), &record); 484 485 assert_eq!( 486 record.transition(&JournalTransition::signed( 487 instance(2), 488 record.revision(), 489 event_id("a"), 490 )), 491 Err(Error::OperationIdentityMismatch) 492 ); 493 assert_eq!( 494 record.transition(&JournalTransition::committed( 495 instance_id, 496 record.revision(), 497 event_id("a"), 498 100, 499 )), 500 Err(Error::InvalidJournalTransition) 501 ); 502 assert_eq!( 503 record.transition(&JournalTransition::cancelled( 504 instance_id, 505 record.revision(), 506 99, 507 )), 508 Err(Error::InvalidJournalTransition) 509 ); 510 let signed = record 511 .transition(&JournalTransition::signed( 512 instance_id, 513 record.revision(), 514 event_id("a"), 515 )) 516 .unwrap(); 517 assert_eq!( 518 signed.transition(&JournalTransition::committed( 519 instance_id, 520 signed.revision(), 521 event_id("b"), 522 101, 523 )), 524 Err(Error::InvalidJournalTransition) 525 ); 526 let recoverable = signed 527 .transition(&JournalTransition::recoverable( 528 instance_id, 529 signed.revision(), 530 recovery, 531 )) 532 .unwrap(); 533 let resumed = recoverable 534 .transition(&JournalTransition::resume( 535 instance_id, 536 recoverable.revision(), 537 )) 538 .unwrap(); 539 assert!(matches!(resumed.state(), JournalState::Signed { .. })); 540 assert_eq!( 541 JournalTransition::resume(instance_id, resumed.revision()).instance_id(), 542 instance_id 543 ); 544 545 assert_eq!( 546 PrepareOperation::new( 547 instance(9), 548 OperationId::SyncPush, 549 key(9), 550 IdempotencyDigest::new([9; 32]), 551 0, 552 ), 553 Err(Error::InvalidOperationTimestamp) 554 ); 555 556 let prepared = prepare(instance(10), 10, 100).into_record().unwrap(); 557 let signed = prepared 558 .transition(&JournalTransition::signed( 559 instance(10), 560 prepared.revision(), 561 event_id("c"), 562 )) 563 .unwrap(); 564 assert_eq!( 565 signed.transition(&JournalTransition::committed( 566 instance(10), 567 signed.revision(), 568 event_id("c"), 569 99, 570 )), 571 Err(Error::InvalidJournalTransition) 572 ); 573 let cancelled_signed = signed 574 .transition(&JournalTransition::cancelled( 575 instance(10), 576 signed.revision(), 577 100, 578 )) 579 .unwrap(); 580 assert_eq!( 581 cancelled_signed.cancellation(), 582 CancellationState::CancelledBeforeCommit 583 ); 584 585 let cancellation_recovery = RecoveryRecord::new( 586 RecoveryPoint::Prepared, 587 RecoveryReason::CancelledBeforeCommit, 588 1, 589 None, 590 ) 591 .unwrap(); 592 let recovered = prepared 593 .transition(&JournalTransition::recoverable( 594 instance(10), 595 prepared.revision(), 596 cancellation_recovery, 597 )) 598 .unwrap(); 599 assert_eq!( 600 recovered.cancellation(), 601 CancellationState::CancelledBeforeCommit 602 ); 603 }