mod.rs (25709B)
1 use crate::SqliteStorage; 2 use crate::backend::map_backend; 3 use radroots_storage::{ 4 Error, Journal, 5 journal::{ 6 BoxFuture, CancellationState, EventId, IdempotencyDigest, IdempotencyKey, JournalRevision, 7 JournalState, JournalTransition, OperationId, OperationInstanceId, OperationRecord, 8 PrepareDisposition, PrepareOperation, PrepareReceipt, RECOVERABLE_QUERY_LIMIT_MAX, 9 RecoveryPoint, RecoveryReason, RecoveryRecord, 10 }, 11 }; 12 use sqlx::{Row, Sqlite}; 13 14 #[cfg_attr(coverage_nightly, coverage(off))] 15 impl Journal for SqliteStorage { 16 fn prepare(&self, operation: PrepareOperation) -> BoxFuture<'_, Result<PrepareReceipt, Error>> { 17 Box::pin(async move { 18 self.require_journal_writer()?; 19 let mut transaction = self 20 .pool() 21 .begin_with("BEGIN IMMEDIATE") 22 .await 23 .map_err(map_backend)?; 24 let receipt = prepare_transaction(&mut transaction, operation).await?; 25 transaction.commit().await.map_err(map_backend)?; 26 Ok(receipt) 27 }) 28 } 29 30 fn operation( 31 &self, 32 instance_id: OperationInstanceId, 33 ) -> BoxFuture<'_, Result<Option<OperationRecord>, Error>> { 34 Box::pin(async move { 35 sqlx::query("SELECT * FROM radroots_runtime_journal_operations WHERE instance_id = ?") 36 .bind(instance_id.as_bytes().as_slice()) 37 .fetch_optional(self.pool()) 38 .await 39 .map_err(map_backend)? 40 .as_ref() 41 .map(decode_record) 42 .transpose() 43 }) 44 } 45 46 fn by_idempotency_key( 47 &self, 48 operation_id: OperationId, 49 idempotency_key: IdempotencyKey, 50 ) -> BoxFuture<'_, Result<Option<OperationRecord>, Error>> { 51 Box::pin(async move { 52 sqlx::query( 53 "SELECT * FROM radroots_runtime_journal_operations 54 WHERE operation_id = ? AND idempotency_key = ?", 55 ) 56 .bind(operation_id.as_str().as_bytes()) 57 .bind(idempotency_key.as_str()) 58 .fetch_optional(self.pool()) 59 .await 60 .map_err(map_backend)? 61 .as_ref() 62 .map(decode_record) 63 .transpose() 64 }) 65 } 66 67 fn transition( 68 &self, 69 transition: JournalTransition, 70 ) -> BoxFuture<'_, Result<OperationRecord, Error>> { 71 Box::pin(async move { 72 self.require_journal_writer()?; 73 let mut transaction = self 74 .pool() 75 .begin_with("BEGIN IMMEDIATE") 76 .await 77 .map_err(map_backend)?; 78 let next = transition_transaction(&mut transaction, transition).await?; 79 transaction.commit().await.map_err(map_backend)?; 80 Ok(next) 81 }) 82 } 83 84 fn recoverable(&self, limit: u16) -> BoxFuture<'_, Result<Vec<OperationRecord>, Error>> { 85 Box::pin(async move { 86 if limit == 0 || limit > RECOVERABLE_QUERY_LIMIT_MAX { 87 return Err(Error::InvalidJournalQueryLimit); 88 } 89 sqlx::query( 90 "SELECT * FROM radroots_runtime_journal_operations 91 WHERE stage = 'recoverable' 92 ORDER BY updated_at_unix_ms, instance_id LIMIT ?", 93 ) 94 .bind(i64::from(limit)) 95 .fetch_all(self.pool()) 96 .await 97 .map_err(map_backend)? 98 .iter() 99 .map(decode_record) 100 .collect() 101 }) 102 } 103 } 104 105 #[cfg_attr(coverage_nightly, coverage(off))] 106 pub(crate) async fn prepare_transaction( 107 transaction: &mut sqlx::Transaction<'_, Sqlite>, 108 operation: PrepareOperation, 109 ) -> Result<PrepareReceipt, Error> { 110 if let Some(row) = sqlx::query( 111 "SELECT * FROM radroots_runtime_journal_operations 112 WHERE idempotency_key = ?", 113 ) 114 .bind(operation.idempotency_key().as_str()) 115 .fetch_optional(&mut **transaction) 116 .await 117 .map_err(map_backend)? 118 { 119 let record = decode_record(&row)?; 120 if record.operation_id() != operation.operation_id() 121 || record.input_digest() != operation.input_digest() 122 || record.instance_id() != operation.instance_id() 123 { 124 return Err(Error::IdempotencyConflict); 125 } 126 return Ok(PrepareReceipt::new(PrepareDisposition::Replay, record)); 127 } 128 if sqlx::query_scalar::<_, i64>( 129 "SELECT 1 FROM radroots_runtime_journal_operations WHERE instance_id = ?", 130 ) 131 .bind(operation.instance_id().as_bytes().as_slice()) 132 .fetch_optional(&mut **transaction) 133 .await 134 .map_err(map_backend)? 135 .is_some() 136 { 137 return Err(Error::OperationIdentityMismatch); 138 } 139 140 let record = operation.into_record()?; 141 insert_record(transaction, &record).await?; 142 Ok(PrepareReceipt::new(PrepareDisposition::Created, record)) 143 } 144 145 #[cfg_attr(coverage_nightly, coverage(off))] 146 pub(crate) async fn transition_transaction( 147 transaction: &mut sqlx::Transaction<'_, Sqlite>, 148 transition: JournalTransition, 149 ) -> Result<OperationRecord, Error> { 150 let row = 151 sqlx::query("SELECT * FROM radroots_runtime_journal_operations WHERE instance_id = ?") 152 .bind(transition.instance_id().as_bytes().as_slice()) 153 .fetch_optional(&mut **transaction) 154 .await 155 .map_err(map_backend)? 156 .ok_or(Error::OperationNotFound)?; 157 let current = decode_record(&row)?; 158 let next = current.transition(&transition)?; 159 let (stage, event_id, recovery, committed_at) = encode_state(next.state()); 160 let result = sqlx::query( 161 "UPDATE radroots_runtime_journal_operations SET 162 revision = ?, stage = ?, event_id = ?, recovery_record = ?, 163 cancellation_state = ?, committed_at_unix_ms = ?, updated_at_unix_ms = ? 164 WHERE instance_id = ? AND revision = ?", 165 ) 166 .bind(i64_from_u64(next.revision().get())?) 167 .bind(stage) 168 .bind(event_id) 169 .bind(recovery) 170 .bind(cancellation_name(next.cancellation())) 171 .bind(committed_at.map(i64_from_u64).transpose()?) 172 .bind(i64_from_u64(updated_at(&next))?) 173 .bind(next.instance_id().as_bytes().as_slice()) 174 .bind(i64_from_u64(current.revision().get())?) 175 .execute(&mut **transaction) 176 .await 177 .map_err(map_backend)?; 178 if result.rows_affected() != 1 { 179 return Err(Error::JournalRevisionConflict); 180 } 181 Ok(next) 182 } 183 184 impl SqliteStorage { 185 fn require_journal_writer(&self) -> Result<(), Error> { 186 if self.event_mode() == radroots_storage::status::EventStoreMode::ReadOnly { 187 return Err(Error::BackendUnavailable); 188 } 189 Ok(()) 190 } 191 } 192 193 #[cfg_attr(coverage_nightly, coverage(off))] 194 async fn insert_record( 195 transaction: &mut sqlx::Transaction<'_, Sqlite>, 196 record: &OperationRecord, 197 ) -> Result<(), Error> { 198 let (stage, event_id, recovery, committed_at) = encode_state(record.state()); 199 sqlx::query( 200 "INSERT INTO radroots_runtime_journal_operations ( 201 instance_id, operation_id, idempotency_key, input_digest, 202 prepared_at_unix_ms, revision, stage, event_id, recovery_record, 203 cancellation_state, committed_at_unix_ms, updated_at_unix_ms 204 ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)", 205 ) 206 .bind(record.instance_id().as_bytes().as_slice()) 207 .bind(record.operation_id().as_str().as_bytes()) 208 .bind(record.idempotency_key().as_str()) 209 .bind(record.input_digest().as_bytes().as_slice()) 210 .bind(i64_from_u64(record.prepared_at_unix_ms())?) 211 .bind(i64_from_u64(record.revision().get())?) 212 .bind(stage) 213 .bind(event_id) 214 .bind(recovery) 215 .bind(cancellation_name(record.cancellation())) 216 .bind(committed_at.map(i64_from_u64).transpose()?) 217 .bind(i64_from_u64(updated_at(record))?) 218 .execute(&mut **transaction) 219 .await 220 .map_err(map_backend)?; 221 Ok(()) 222 } 223 224 pub(crate) fn decode_record(row: &sqlx::sqlite::SqliteRow) -> Result<OperationRecord, Error> { 225 let instance_id = OperationInstanceId::new(array( 226 row.try_get::<Vec<u8>, _>("instance_id") 227 .map_err(map_corrupt)?, 228 )?) 229 .map_err(|_| Error::CorruptJournalRecord)?; 230 let operation_text = String::from_utf8( 231 row.try_get::<Vec<u8>, _>("operation_id") 232 .map_err(map_corrupt)?, 233 ) 234 .map_err(|_| Error::CorruptJournalRecord)?; 235 let operation_id = 236 OperationId::parse(operation_text.as_str()).map_err(|_| Error::CorruptJournalRecord)?; 237 let idempotency_key = IdempotencyKey::parse( 238 row.try_get::<String, _>("idempotency_key") 239 .map_err(map_corrupt)?, 240 ) 241 .map_err(|_| Error::CorruptJournalRecord)?; 242 let input_digest = IdempotencyDigest::new(array( 243 row.try_get::<Vec<u8>, _>("input_digest") 244 .map_err(map_corrupt)?, 245 )?); 246 let prepared_at = u64_from_i64(row.try_get("prepared_at_unix_ms").map_err(map_corrupt)?)?; 247 let revision = 248 JournalRevision::new(u64_from_i64(row.try_get("revision").map_err(map_corrupt)?)?) 249 .map_err(|_| Error::CorruptJournalRecord)?; 250 let event_id = row 251 .try_get::<Option<Vec<u8>>, _>("event_id") 252 .map_err(map_corrupt)? 253 .map(|bytes| array(bytes).map(EventId::from_bytes)) 254 .transpose()?; 255 let recovery = row 256 .try_get::<Option<Vec<u8>>, _>("recovery_record") 257 .map_err(map_corrupt)? 258 .map(|bytes| decode_recovery(bytes.as_slice())) 259 .transpose()?; 260 let committed_at = row 261 .try_get::<Option<i64>, _>("committed_at_unix_ms") 262 .map_err(map_corrupt)? 263 .map(u64_from_i64) 264 .transpose()?; 265 let state = decode_state( 266 row.try_get::<String, _>("stage") 267 .map_err(map_corrupt)? 268 .as_str(), 269 event_id, 270 recovery, 271 committed_at, 272 )?; 273 let cancellation = cancellation( 274 row.try_get::<String, _>("cancellation_state") 275 .map_err(map_corrupt)? 276 .as_str(), 277 )?; 278 OperationRecord::from_parts( 279 instance_id, 280 operation_id, 281 idempotency_key, 282 input_digest, 283 prepared_at, 284 revision, 285 state, 286 cancellation, 287 ) 288 .map_err(|_| Error::CorruptJournalRecord) 289 } 290 291 type EncodedState = (&'static str, Option<Vec<u8>>, Option<Vec<u8>>, Option<u64>); 292 293 fn encode_state(state: &JournalState) -> EncodedState { 294 match state { 295 JournalState::Prepared => ("prepared", None, None, None), 296 JournalState::Signed { event_id } => { 297 ("signed", Some(event_id.as_bytes().to_vec()), None, None) 298 } 299 JournalState::Recoverable(record) => { 300 let event_id = match record.point() { 301 RecoveryPoint::Prepared => None, 302 RecoveryPoint::Signed { event_id } => Some(event_id.as_bytes().to_vec()), 303 }; 304 ("recoverable", event_id, Some(encode_recovery(record)), None) 305 } 306 JournalState::Committed { 307 event_id, 308 committed_at_unix_ms, 309 } => ( 310 "committed", 311 Some(event_id.as_bytes().to_vec()), 312 None, 313 Some(*committed_at_unix_ms), 314 ), 315 } 316 } 317 318 fn decode_state( 319 stage: &str, 320 event_id: Option<EventId>, 321 recovery: Option<RecoveryRecord>, 322 committed_at: Option<u64>, 323 ) -> Result<JournalState, Error> { 324 match (stage, event_id, recovery, committed_at) { 325 ("prepared", None, None, None) => Ok(JournalState::Prepared), 326 ("signed", Some(event_id), None, None) => Ok(JournalState::Signed { event_id }), 327 ("recoverable", event_id, Some(recovery), None) 328 if recovery_event_id(&recovery) == event_id => 329 { 330 Ok(JournalState::Recoverable(recovery)) 331 } 332 ("committed", Some(event_id), None, Some(committed_at_unix_ms)) => { 333 Ok(JournalState::Committed { 334 event_id, 335 committed_at_unix_ms, 336 }) 337 } 338 _ => Err(Error::CorruptJournalRecord), 339 } 340 } 341 342 fn recovery_event_id(record: &RecoveryRecord) -> Option<EventId> { 343 match record.point() { 344 RecoveryPoint::Prepared => None, 345 RecoveryPoint::Signed { event_id } => Some(*event_id), 346 } 347 } 348 349 fn encode_recovery(record: &RecoveryRecord) -> Vec<u8> { 350 let mut bytes = Vec::with_capacity(48); 351 bytes.push(1); 352 match record.point() { 353 RecoveryPoint::Prepared => bytes.push(0), 354 RecoveryPoint::Signed { event_id } => { 355 bytes.push(1); 356 bytes.extend_from_slice(event_id.as_bytes()); 357 } 358 } 359 bytes.push(recovery_reason_byte(record.reason())); 360 bytes.extend_from_slice(&record.attempt().to_be_bytes()); 361 match record.retry_not_before_unix_ms() { 362 Some(value) => { 363 bytes.push(1); 364 bytes.extend_from_slice(&value.to_be_bytes()); 365 } 366 None => bytes.push(0), 367 } 368 bytes 369 } 370 371 fn decode_recovery(bytes: &[u8]) -> Result<RecoveryRecord, Error> { 372 let mut offset = 0; 373 if take_byte(bytes, &mut offset)? != 1 { 374 return Err(Error::CorruptJournalRecord); 375 } 376 let point = match take_byte(bytes, &mut offset)? { 377 0 => RecoveryPoint::Prepared, 378 1 => RecoveryPoint::Signed { 379 event_id: EventId::from_bytes(take_array(bytes, &mut offset)?), 380 }, 381 _ => return Err(Error::CorruptJournalRecord), 382 }; 383 let reason = recovery_reason(take_byte(bytes, &mut offset)?)?; 384 let attempt = u32::from_be_bytes(take_array(bytes, &mut offset)?); 385 let retry = match take_byte(bytes, &mut offset)? { 386 0 => None, 387 1 => Some(u64::from_be_bytes(take_array(bytes, &mut offset)?)), 388 _ => return Err(Error::CorruptJournalRecord), 389 }; 390 if offset != bytes.len() { 391 return Err(Error::CorruptJournalRecord); 392 } 393 RecoveryRecord::new(point, reason, attempt, retry).map_err(|_| Error::CorruptJournalRecord) 394 } 395 396 const fn recovery_reason_byte(reason: RecoveryReason) -> u8 { 397 match reason { 398 RecoveryReason::CancelledBeforeCommit => 0, 399 RecoveryReason::SignerUnavailable => 1, 400 RecoveryReason::TransportUnavailable => 2, 401 RecoveryReason::StorageUnavailable => 3, 402 RecoveryReason::DeadlineExceeded => 4, 403 RecoveryReason::Interrupted => 5, 404 } 405 } 406 407 const fn recovery_reason(value: u8) -> Result<RecoveryReason, Error> { 408 match value { 409 0 => Ok(RecoveryReason::CancelledBeforeCommit), 410 1 => Ok(RecoveryReason::SignerUnavailable), 411 2 => Ok(RecoveryReason::TransportUnavailable), 412 3 => Ok(RecoveryReason::StorageUnavailable), 413 4 => Ok(RecoveryReason::DeadlineExceeded), 414 5 => Ok(RecoveryReason::Interrupted), 415 _ => Err(Error::CorruptJournalRecord), 416 } 417 } 418 419 const fn cancellation_name(value: CancellationState) -> &'static str { 420 match value { 421 CancellationState::NotRequested => "not_requested", 422 CancellationState::CancelledBeforeCommit => "cancelled_before_commit", 423 CancellationState::ObservedAfterCommit => "observed_after_commit", 424 } 425 } 426 427 const fn cancellation(value: &str) -> Result<CancellationState, Error> { 428 match value.as_bytes() { 429 b"not_requested" => Ok(CancellationState::NotRequested), 430 b"cancelled_before_commit" => Ok(CancellationState::CancelledBeforeCommit), 431 b"observed_after_commit" => Ok(CancellationState::ObservedAfterCommit), 432 _ => Err(Error::CorruptJournalRecord), 433 } 434 } 435 436 fn updated_at(record: &OperationRecord) -> u64 { 437 match record.state() { 438 JournalState::Committed { 439 committed_at_unix_ms, 440 .. 441 } => *committed_at_unix_ms, 442 JournalState::Recoverable(recovery) => recovery 443 .retry_not_before_unix_ms() 444 .unwrap_or(record.prepared_at_unix_ms()), 445 JournalState::Prepared | JournalState::Signed { .. } => record.prepared_at_unix_ms(), 446 } 447 } 448 449 fn take_byte(bytes: &[u8], offset: &mut usize) -> Result<u8, Error> { 450 let value = bytes 451 .get(*offset) 452 .copied() 453 .ok_or(Error::CorruptJournalRecord)?; 454 *offset += 1; 455 Ok(value) 456 } 457 458 fn take_array<const N: usize>(bytes: &[u8], offset: &mut usize) -> Result<[u8; N], Error> { 459 let end = offset.checked_add(N).ok_or(Error::CorruptJournalRecord)?; 460 let value = bytes 461 .get(*offset..end) 462 .ok_or(Error::CorruptJournalRecord)? 463 .try_into() 464 .map_err(|_| Error::CorruptJournalRecord)?; 465 *offset = end; 466 Ok(value) 467 } 468 469 fn array<const N: usize>(bytes: Vec<u8>) -> Result<[u8; N], Error> { 470 bytes.try_into().map_err(|_| Error::CorruptJournalRecord) 471 } 472 473 fn i64_from_u64(value: u64) -> Result<i64, Error> { 474 i64::try_from(value).map_err(|_| Error::CorruptJournalRecord) 475 } 476 477 fn u64_from_i64(value: i64) -> Result<u64, Error> { 478 u64::try_from(value).map_err(|_| Error::CorruptJournalRecord) 479 } 480 481 fn map_corrupt(_: sqlx::Error) -> Error { 482 Error::CorruptJournalRecord 483 } 484 485 #[cfg(test)] 486 #[cfg_attr(coverage_nightly, coverage(off))] 487 mod tests { 488 use super::*; 489 use crate::migration::runtime::{MIGRATIONS, migration_sql}; 490 use radroots_storage::{ 491 Journal, event::SourceGeneration, journal::JournalStage, status::EventStoreMode, 492 }; 493 use sqlx::sqlite::SqlitePoolOptions; 494 495 async fn store(mode: EventStoreMode) -> SqliteStorage { 496 let pool = SqlitePoolOptions::new() 497 .max_connections(1) 498 .connect("sqlite::memory:") 499 .await 500 .expect("memory SQLite"); 501 sqlx::query("PRAGMA foreign_keys = ON") 502 .execute(&pool) 503 .await 504 .expect("foreign keys"); 505 for migration in MIGRATIONS { 506 sqlx::raw_sql(migration_sql(migration.version()).expect("registered SQL")) 507 .execute(&pool) 508 .await 509 .expect("runtime migration"); 510 } 511 SqliteStorage::new( 512 pool, 513 SourceGeneration::new([21; 32]).expect("generation"), 514 mode, 515 ) 516 } 517 518 fn instance(byte: u8) -> OperationInstanceId { 519 OperationInstanceId::new([byte; 16]).expect("instance") 520 } 521 522 fn key(byte: u8) -> IdempotencyKey { 523 IdempotencyKey::parse(format!("journal-{byte:02x}")).expect("key") 524 } 525 526 fn prepare( 527 instance_id: OperationInstanceId, 528 key_byte: u8, 529 digest: u8, 530 at: u64, 531 ) -> PrepareOperation { 532 PrepareOperation::new( 533 instance_id, 534 OperationId::SyncPush, 535 key(key_byte), 536 IdempotencyDigest::new([digest; 32]), 537 at, 538 ) 539 .expect("prepare") 540 } 541 542 #[test] 543 fn recovery_decoder_and_state_guard_reject_noncanonical_encodings() { 544 let event_id = EventId::from_bytes([42; 32]); 545 let recovery = RecoveryRecord::new( 546 RecoveryPoint::Signed { event_id }, 547 RecoveryReason::Interrupted, 548 1, 549 None, 550 ) 551 .expect("recovery"); 552 assert_eq!( 553 decode_state("recoverable", None, Some(recovery.clone()), None), 554 Err(Error::CorruptJournalRecord) 555 ); 556 assert_eq!(decode_recovery(&[2]), Err(Error::CorruptJournalRecord)); 557 let mut encoded = encode_recovery(&recovery); 558 encoded.push(0); 559 assert_eq!( 560 decode_recovery(encoded.as_slice()), 561 Err(Error::CorruptJournalRecord) 562 ); 563 } 564 565 #[tokio::test] 566 async fn prepare_replays_exact_identity_and_rejects_conflicts() { 567 let store = store(EventStoreMode::ReadWrite).await; 568 let request = prepare(instance(1), 1, 2, 100); 569 let created = store.prepare(request.clone()).await.expect("created"); 570 assert_eq!(created.disposition(), PrepareDisposition::Created); 571 assert_eq!(created.record().revision(), JournalRevision::INITIAL); 572 let replay = store.prepare(request).await.expect("replay"); 573 assert_eq!(replay.disposition(), PrepareDisposition::Replay); 574 assert_eq!(replay.record(), created.record()); 575 assert_eq!( 576 store.prepare(prepare(instance(1), 1, 3, 100)).await, 577 Err(Error::IdempotencyConflict) 578 ); 579 assert_eq!( 580 store.prepare(prepare(instance(1), 2, 2, 100)).await, 581 Err(Error::OperationIdentityMismatch) 582 ); 583 assert_eq!( 584 store 585 .operation(instance(1)) 586 .await 587 .expect("lookup") 588 .expect("record"), 589 *created.record() 590 ); 591 assert_eq!( 592 store 593 .by_idempotency_key(OperationId::SyncPush, key(1)) 594 .await 595 .expect("key lookup") 596 .expect("record"), 597 *created.record() 598 ); 599 assert!( 600 store 601 .by_idempotency_key(OperationId::FarmPublish, key(1)) 602 .await 603 .expect("wrong operation lookup") 604 .is_none() 605 ); 606 } 607 608 #[tokio::test] 609 async fn lifecycle_recovery_commit_and_cancellation_round_trip() { 610 let store = store(EventStoreMode::ReadWrite).await; 611 let instance_id = instance(3); 612 let event_id = EventId::from_bytes([4; 32]); 613 let prepared = store 614 .prepare(prepare(instance_id, 3, 3, 100)) 615 .await 616 .expect("prepare") 617 .record() 618 .clone(); 619 let signed = store 620 .transition(JournalTransition::signed( 621 instance_id, 622 prepared.revision(), 623 event_id, 624 )) 625 .await 626 .expect("signed"); 627 assert_eq!(signed.state().stage(), JournalStage::Signed); 628 assert_eq!( 629 store 630 .transition(JournalTransition::signed( 631 instance_id, 632 prepared.revision(), 633 event_id, 634 )) 635 .await, 636 Err(Error::JournalRevisionConflict) 637 ); 638 639 let recovery = RecoveryRecord::new( 640 RecoveryPoint::Signed { event_id }, 641 RecoveryReason::TransportUnavailable, 642 2, 643 Some(200), 644 ) 645 .expect("recovery"); 646 let recoverable = store 647 .transition(JournalTransition::recoverable( 648 instance_id, 649 signed.revision(), 650 recovery.clone(), 651 )) 652 .await 653 .expect("recoverable"); 654 assert_eq!( 655 store.recoverable(10).await.expect("recovery query"), 656 vec![recoverable.clone()] 657 ); 658 assert_eq!(recoverable.state(), &JournalState::Recoverable(recovery)); 659 660 let resumed = store 661 .transition(JournalTransition::resume( 662 instance_id, 663 recoverable.revision(), 664 )) 665 .await 666 .expect("resume"); 667 assert_eq!(resumed.state().stage(), JournalStage::Signed); 668 let committed = store 669 .transition(JournalTransition::committed( 670 instance_id, 671 resumed.revision(), 672 event_id, 673 250, 674 )) 675 .await 676 .expect("commit"); 677 assert_eq!(committed.state().stage(), JournalStage::Committed); 678 let cancelled = store 679 .transition(JournalTransition::cancelled( 680 instance_id, 681 committed.revision(), 682 260, 683 )) 684 .await 685 .expect("post-commit cancellation"); 686 assert_eq!(cancelled.state(), committed.state()); 687 assert_eq!( 688 cancelled.cancellation(), 689 CancellationState::ObservedAfterCommit 690 ); 691 assert!( 692 store 693 .recoverable(10) 694 .await 695 .expect("empty recovery") 696 .is_empty() 697 ); 698 } 699 700 #[tokio::test] 701 async fn cancellation_corruption_bounds_and_read_only_mode_fail_closed() { 702 let store = store(EventStoreMode::ReadWrite).await; 703 let instance_id = instance(5); 704 let prepared = store 705 .prepare(prepare(instance_id, 5, 5, 500)) 706 .await 707 .expect("prepare") 708 .record() 709 .clone(); 710 let cancelled = store 711 .transition(JournalTransition::cancelled( 712 instance_id, 713 prepared.revision(), 714 501, 715 )) 716 .await 717 .expect("cancel"); 718 assert_eq!(cancelled.state().stage(), JournalStage::Recoverable); 719 assert_eq!( 720 store.recoverable(0).await, 721 Err(Error::InvalidJournalQueryLimit) 722 ); 723 assert_eq!( 724 store.recoverable(RECOVERABLE_QUERY_LIMIT_MAX + 1).await, 725 Err(Error::InvalidJournalQueryLimit) 726 ); 727 728 sqlx::query( 729 "UPDATE radroots_runtime_journal_operations 730 SET recovery_record = X'0100FF' WHERE instance_id = ?", 731 ) 732 .bind(instance_id.as_bytes().as_slice()) 733 .execute(store.pool()) 734 .await 735 .expect("forge corrupt recovery"); 736 assert_eq!( 737 store.operation(instance_id).await, 738 Err(Error::CorruptJournalRecord) 739 ); 740 741 let read_only = SqliteStorage::new( 742 store.pool().clone(), 743 SourceGeneration::new([21; 32]).expect("generation"), 744 EventStoreMode::ReadOnly, 745 ); 746 assert_eq!( 747 read_only.prepare(prepare(instance(6), 6, 6, 600)).await, 748 Err(Error::BackendUnavailable) 749 ); 750 } 751 }