atomic.rs (19148B)
1 use crate::backend::map_backend; 2 use crate::{SqliteStorage, projection}; 3 use radroots_storage::{ 4 Error, 5 atomic::{ 6 AtomicCommit, AtomicCommitDigest, AtomicCommitDisposition, AtomicCommitId, 7 AtomicCommitOutcome, AtomicCommitReceipt, AtomicStorage, AtomicWorkflow, 8 AtomicWorkflowKind, 9 }, 10 event::{ 11 AdmissionDisposition, AdmissionReceipt, AdmissionStage, BoxFuture, EventPosition, 12 EventSequence, SourceGeneration, 13 }, 14 }; 15 use sqlx::{Row, Sqlite}; 16 17 const RECEIPT_FORMAT_VERSION: u8 = 1; 18 const RECEIPT_MAX_BYTES: usize = 4 * 1024 * 1024; 19 20 #[cfg_attr(coverage_nightly, coverage(off))] 21 impl AtomicStorage for SqliteStorage { 22 fn commit(&self, request: AtomicCommit) -> BoxFuture<'_, Result<AtomicCommitReceipt, Error>> { 23 Box::pin(async move { 24 self.require_atomic_writer()?; 25 let mut transaction = self 26 .pool() 27 .begin_with("BEGIN IMMEDIATE") 28 .await 29 .map_err(map_backend)?; 30 31 let result = commit_transaction(self, &mut transaction, &request).await; 32 match result { 33 Ok(receipt) => { 34 transaction.commit().await.map_err(map_backend)?; 35 Ok(receipt) 36 } 37 Err(primary) => { 38 let rollback = transaction.rollback().await; 39 Err(preserve_primary(primary, rollback)) 40 } 41 } 42 }) 43 } 44 45 fn receipt( 46 &self, 47 commit_id: AtomicCommitId, 48 ) -> BoxFuture<'_, Result<Option<AtomicCommitReceipt>, Error>> { 49 Box::pin(async move { 50 sqlx::query( 51 "SELECT commit_id, commit_digest, workflow_kind, requested_at_unix_ms, 52 committed_at_unix_ms, receipt 53 FROM radroots_runtime_atomic_commits WHERE commit_id = ?", 54 ) 55 .bind(commit_id.as_bytes().as_slice()) 56 .fetch_optional(self.pool()) 57 .await 58 .map_err(map_backend)? 59 .as_ref() 60 .map(decode_receipt_row) 61 .transpose() 62 }) 63 } 64 } 65 66 impl SqliteStorage { 67 fn require_atomic_writer(&self) -> Result<(), Error> { 68 if self.event_mode() == radroots_storage::status::EventStoreMode::ReadOnly { 69 return Err(Error::BackendUnavailable); 70 } 71 Ok(()) 72 } 73 } 74 75 async fn commit_transaction( 76 storage: &SqliteStorage, 77 transaction: &mut sqlx::Transaction<'_, Sqlite>, 78 request: &AtomicCommit, 79 ) -> Result<AtomicCommitReceipt, Error> { 80 if let Some(row) = sqlx::query( 81 "SELECT commit_id, commit_digest, workflow_kind, requested_at_unix_ms, 82 committed_at_unix_ms, receipt 83 FROM radroots_runtime_atomic_commits WHERE commit_id = ?", 84 ) 85 .bind(request.commit_id().as_bytes().as_slice()) 86 .fetch_optional(&mut **transaction) 87 .await 88 .map_err(map_backend)? 89 { 90 let committed = decode_receipt_row(&row)?; 91 if [ 92 committed.digest() != request.digest(), 93 committed.outcome().kind() != request.workflow().kind(), 94 ] 95 .contains(&true) 96 { 97 return Err(Error::AtomicCommitConflict); 98 } 99 return AtomicCommitReceipt::new( 100 request, 101 AtomicCommitDisposition::Replay, 102 committed.committed_at_unix_ms(), 103 committed.outcome().clone(), 104 ); 105 } 106 107 let outcome = execute_workflow(storage, transaction, request.workflow().clone()).await?; 108 let committed_at_unix_ms = request.requested_at_unix_ms(); 109 let receipt = AtomicCommitReceipt::new( 110 request, 111 AtomicCommitDisposition::Committed, 112 committed_at_unix_ms, 113 outcome, 114 )?; 115 let snapshot = encode_outcome(receipt.outcome())?; 116 sqlx::query( 117 "INSERT INTO radroots_runtime_atomic_commits ( 118 commit_id, commit_digest, workflow_kind, requested_at_unix_ms, 119 committed_at_unix_ms, receipt 120 ) VALUES (?, ?, ?, ?, ?, ?)", 121 ) 122 .bind(request.commit_id().as_bytes().as_slice()) 123 .bind(request.digest().as_bytes().as_slice()) 124 .bind(workflow_name(request.workflow().kind())) 125 .bind(i64_from_u64(request.requested_at_unix_ms())?) 126 .bind(i64_from_u64(committed_at_unix_ms)?) 127 .bind(snapshot) 128 .execute(&mut **transaction) 129 .await 130 .map_err(map_backend)?; 131 Ok(receipt) 132 } 133 134 async fn execute_workflow( 135 storage: &SqliteStorage, 136 transaction: &mut sqlx::Transaction<'_, Sqlite>, 137 workflow: AtomicWorkflow, 138 ) -> Result<AtomicCommitOutcome, Error> { 139 match workflow { 140 AtomicWorkflow::Ingested(ingested) => { 141 let admission = storage 142 .admit_transaction(transaction, ingested.admission().clone()) 143 .await?; 144 let projection = match ingested.projection().cloned() { 145 Some(checkpoint) => Some(Box::new( 146 projection::checkpoint_transaction(transaction, checkpoint).await?, 147 )), 148 None => None, 149 }; 150 Ok(AtomicCommitOutcome::Ingested { 151 admission, 152 projection, 153 }) 154 } 155 } 156 } 157 158 fn decode_receipt_row(row: &sqlx::sqlite::SqliteRow) -> Result<AtomicCommitReceipt, Error> { 159 let commit_id = AtomicCommitId::new(array( 160 row.try_get::<Vec<u8>, _>("commit_id") 161 .map_err(map_corrupt)?, 162 )?) 163 .map_err(|_| Error::AtomicCommitFailed)?; 164 let digest = AtomicCommitDigest::new(array( 165 row.try_get::<Vec<u8>, _>("commit_digest") 166 .map_err(map_corrupt)?, 167 )?); 168 let workflow_kind = workflow_kind( 169 row.try_get::<String, _>("workflow_kind") 170 .map_err(map_corrupt)? 171 .as_str(), 172 )?; 173 let requested_at_unix_ms = 174 u64_from_i64(row.try_get("requested_at_unix_ms").map_err(map_corrupt)?)?; 175 let committed_at_unix_ms = 176 u64_from_i64(row.try_get("committed_at_unix_ms").map_err(map_corrupt)?)?; 177 let bytes = row.try_get::<Vec<u8>, _>("receipt").map_err(map_corrupt)?; 178 let outcome = decode_outcome(bytes.as_slice())?; 179 AtomicCommitReceipt::from_durable_parts( 180 commit_id, 181 digest, 182 AtomicCommitDisposition::Committed, 183 requested_at_unix_ms, 184 committed_at_unix_ms, 185 workflow_kind, 186 outcome, 187 ) 188 .map_err(|_| Error::AtomicCommitFailed) 189 } 190 191 fn encode_outcome(outcome: &AtomicCommitOutcome) -> Result<Vec<u8>, Error> { 192 let mut bytes = Vec::with_capacity(512); 193 bytes.push(RECEIPT_FORMAT_VERSION); 194 match outcome { 195 AtomicCommitOutcome::Ingested { 196 admission, 197 projection, 198 } => { 199 bytes.push(4); 200 encode_admission(&mut bytes, admission); 201 match projection { 202 Some(status) => { 203 bytes.push(1); 204 put_blob( 205 &mut bytes, 206 projection::encode_status_snapshot(status)?.as_slice(), 207 )?; 208 } 209 None => bytes.push(0), 210 } 211 } 212 } 213 if bytes.len() > RECEIPT_MAX_BYTES { 214 return Err(Error::AtomicCommitFailed); 215 } 216 Ok(bytes) 217 } 218 219 fn decode_outcome(bytes: &[u8]) -> Result<AtomicCommitOutcome, Error> { 220 if bytes.len() > RECEIPT_MAX_BYTES { 221 return Err(Error::AtomicCommitFailed); 222 } 223 let mut cursor = Cursor::new(bytes); 224 if cursor.byte()? != RECEIPT_FORMAT_VERSION { 225 return Err(Error::AtomicCommitFailed); 226 } 227 let outcome = match cursor.byte()? { 228 4 => { 229 let admission = decode_admission(&mut cursor)?; 230 let projection = match cursor.byte()? { 231 0 => None, 232 1 => Some(Box::new( 233 projection::decode_status_snapshot(cursor.blob()?) 234 .map_err(|_| Error::AtomicCommitFailed)?, 235 )), 236 _ => return Err(Error::AtomicCommitFailed), 237 }; 238 AtomicCommitOutcome::Ingested { 239 admission, 240 projection, 241 } 242 } 243 _ => return Err(Error::AtomicCommitFailed), 244 }; 245 cursor.finish()?; 246 Ok(outcome) 247 } 248 249 fn encode_admission(bytes: &mut Vec<u8>, receipt: &AdmissionReceipt) { 250 bytes.extend_from_slice(receipt.event_id().as_bytes()); 251 bytes.extend_from_slice(receipt.position().generation().as_bytes()); 252 bytes.extend_from_slice(&receipt.position().sequence().get().to_be_bytes()); 253 bytes.push(match receipt.stage() { 254 AdmissionStage::Raw => 0, 255 AdmissionStage::Verified => 1, 256 AdmissionStage::Visible => 2, 257 }); 258 bytes.push(match receipt.disposition() { 259 AdmissionDisposition::Inserted => 0, 260 AdmissionDisposition::Advanced => 1, 261 AdmissionDisposition::Duplicate => 2, 262 }); 263 } 264 265 fn decode_admission(cursor: &mut Cursor<'_>) -> Result<AdmissionReceipt, Error> { 266 let event_id = radroots_storage::event::EventId::from_bytes(cursor.array()?); 267 let generation = 268 SourceGeneration::new(cursor.array()?).map_err(|_| Error::AtomicCommitFailed)?; 269 let sequence = EventSequence::new(cursor.u64()?).map_err(|_| Error::AtomicCommitFailed)?; 270 let stage = match cursor.byte()? { 271 0 => AdmissionStage::Raw, 272 1 => AdmissionStage::Verified, 273 2 => AdmissionStage::Visible, 274 _ => return Err(Error::AtomicCommitFailed), 275 }; 276 let disposition = match cursor.byte()? { 277 0 => AdmissionDisposition::Inserted, 278 1 => AdmissionDisposition::Advanced, 279 2 => AdmissionDisposition::Duplicate, 280 _ => return Err(Error::AtomicCommitFailed), 281 }; 282 Ok(AdmissionReceipt::new( 283 event_id, 284 EventPosition::new(generation, sequence), 285 stage, 286 disposition, 287 )) 288 } 289 290 fn put_blob(bytes: &mut Vec<u8>, value: &[u8]) -> Result<(), Error> { 291 let length = u32::try_from(value.len()).map_err(|_| Error::AtomicCommitFailed)?; 292 bytes.extend_from_slice(&length.to_be_bytes()); 293 bytes.extend_from_slice(value); 294 Ok(()) 295 } 296 297 struct Cursor<'a> { 298 bytes: &'a [u8], 299 offset: usize, 300 } 301 302 impl<'a> Cursor<'a> { 303 const fn new(bytes: &'a [u8]) -> Self { 304 Self { bytes, offset: 0 } 305 } 306 307 fn byte(&mut self) -> Result<u8, Error> { 308 let value = self 309 .bytes 310 .get(self.offset) 311 .copied() 312 .ok_or(Error::AtomicCommitFailed)?; 313 self.offset += 1; 314 Ok(value) 315 } 316 317 fn u32(&mut self) -> Result<u32, Error> { 318 Ok(u32::from_be_bytes(self.array()?)) 319 } 320 321 fn u64(&mut self) -> Result<u64, Error> { 322 Ok(u64::from_be_bytes(self.array()?)) 323 } 324 325 fn array<const N: usize>(&mut self) -> Result<[u8; N], Error> { 326 self.take(N)? 327 .try_into() 328 .map_err(|_| Error::AtomicCommitFailed) 329 } 330 331 fn blob(&mut self) -> Result<&'a [u8], Error> { 332 let length = usize::try_from(self.u32()?).map_err(|_| Error::AtomicCommitFailed)?; 333 self.take(length) 334 } 335 336 fn take(&mut self, length: usize) -> Result<&'a [u8], Error> { 337 let end = self 338 .offset 339 .checked_add(length) 340 .ok_or(Error::AtomicCommitFailed)?; 341 let value = self 342 .bytes 343 .get(self.offset..end) 344 .ok_or(Error::AtomicCommitFailed)?; 345 self.offset = end; 346 Ok(value) 347 } 348 349 fn finish(self) -> Result<(), Error> { 350 if self.offset == self.bytes.len() { 351 Ok(()) 352 } else { 353 Err(Error::AtomicCommitFailed) 354 } 355 } 356 } 357 358 fn preserve_primary<T>(primary: Error, _rollback: Result<(), T>) -> Error { 359 primary 360 } 361 362 const fn workflow_name(kind: AtomicWorkflowKind) -> &'static str { 363 match kind { 364 AtomicWorkflowKind::Ingested => "ingested", 365 } 366 } 367 368 fn workflow_kind(value: &str) -> Result<AtomicWorkflowKind, Error> { 369 match value.as_bytes() { 370 b"ingested" => Ok(AtomicWorkflowKind::Ingested), 371 _ => Err(Error::AtomicCommitFailed), 372 } 373 } 374 375 fn array<const N: usize>(bytes: Vec<u8>) -> Result<[u8; N], Error> { 376 bytes.try_into().map_err(|_| Error::AtomicCommitFailed) 377 } 378 379 fn i64_from_u64(value: u64) -> Result<i64, Error> { 380 i64::try_from(value).map_err(|_| Error::AtomicCommitFailed) 381 } 382 383 fn u64_from_i64(value: i64) -> Result<u64, Error> { 384 u64::try_from(value).map_err(|_| Error::AtomicCommitFailed) 385 } 386 387 fn map_corrupt(_: sqlx::Error) -> Error { 388 Error::AtomicCommitFailed 389 } 390 391 #[cfg(test)] 392 #[cfg_attr(coverage_nightly, coverage(off))] 393 mod tests { 394 use super::*; 395 use crate::migration::runtime::{MIGRATIONS, migration_sql}; 396 use radroots_event::{SignedEvent, wire::Nip01EventWire}; 397 use radroots_storage::{ 398 atomic::{AtomicWorkflow, CommitIngested}, 399 event::{EventAdmission, SourceGeneration}, 400 projection::{ProjectionCheckpoint, ProjectionGeneration, ProjectionId}, 401 status::EventStoreMode, 402 }; 403 use radroots_transport::{ 404 Target, TransportId, 405 source::{EventProvenance, ObservedEvent}, 406 }; 407 use sqlx::sqlite::SqlitePoolOptions; 408 409 async fn store(mode: EventStoreMode) -> SqliteStorage { 410 let generation = SourceGeneration::new([71; 32]).expect("generation"); 411 let pool = SqlitePoolOptions::new() 412 .max_connections(1) 413 .connect("sqlite::memory:") 414 .await 415 .expect("memory SQLite"); 416 sqlx::query("PRAGMA foreign_keys = ON") 417 .execute(&pool) 418 .await 419 .expect("foreign keys"); 420 for migration in MIGRATIONS { 421 sqlx::raw_sql(migration_sql(migration.version()).expect("registered SQL")) 422 .execute(&pool) 423 .await 424 .expect("runtime migration"); 425 } 426 sqlx::query( 427 "INSERT INTO radroots_runtime_source_generations ( 428 generation, sequence_head, state, created_at_unix_ms, retired_at_unix_ms 429 ) VALUES (?, 0, 'active', 1, NULL)", 430 ) 431 .bind(generation.as_bytes().as_slice()) 432 .execute(&pool) 433 .await 434 .expect("source generation"); 435 SqliteStorage::new(pool, generation, mode) 436 } 437 438 fn signed_event() -> SignedEvent { 439 let mut wire = Nip01EventWire { 440 id: "0".repeat(64), 441 pubkey: "585591529da0bab31b3b1b1f986611cf5f435dca84f978c89ee8a40cca7103df".to_owned(), 442 created_at: 1_800_001_001, 443 kind: 1, 444 tags: vec![], 445 content: "atomic inbound".to_owned(), 446 sig: "42".repeat(64), 447 extra: Default::default(), 448 }; 449 wire.id = wire.computed_event_id().expect("event id").to_hex(); 450 let raw_json = serde_json::json!({ 451 "id": &wire.id, 452 "pubkey": &wire.pubkey, 453 "created_at": wire.created_at, 454 "kind": wire.kind, 455 "tags": &wire.tags, 456 "content": &wire.content, 457 "sig": &wire.sig, 458 }) 459 .to_string(); 460 SignedEvent::from_wire_verified_id(wire, raw_json).expect("signed event") 461 } 462 463 fn admission(observed_at: u64) -> EventAdmission { 464 let target = Target::new(TransportId::NOSTR, "wss://atomic.example").expect("target"); 465 let provenance = EventProvenance::new( 466 TransportId::NOSTR, 467 target.fingerprint().clone(), 468 observed_at, 469 ) 470 .expect("provenance"); 471 EventAdmission::raw(ObservedEvent::new(signed_event(), provenance)) 472 } 473 474 fn commit(byte: u8, digest: u8, workflow: AtomicWorkflow) -> AtomicCommit { 475 AtomicCommit::new( 476 AtomicCommitId::new([byte; 16]).expect("commit id"), 477 AtomicCommitDigest::new([digest; 32]), 478 100, 479 workflow, 480 ) 481 .expect("atomic commit") 482 } 483 484 fn checkpoint() -> ProjectionCheckpoint { 485 ProjectionCheckpoint::new( 486 ProjectionId::parse("atomic_projection").expect("projection id"), 487 ProjectionGeneration::new([81; 32]).expect("projection generation"), 488 None, 489 1, 490 100, 491 ) 492 .expect("checkpoint") 493 } 494 495 #[tokio::test] 496 async fn inbound_commit_is_atomic_replayable_and_reconstructable() { 497 let store = store(EventStoreMode::ReadWrite).await; 498 let request = commit( 499 1, 500 1, 501 AtomicWorkflow::Ingested(Box::new(CommitIngested::new( 502 admission(90), 503 Some(checkpoint()), 504 ))), 505 ); 506 let committed = store.commit(request.clone()).await.expect("commit"); 507 assert_eq!(committed.disposition(), AtomicCommitDisposition::Committed); 508 assert!(matches!( 509 committed.outcome(), 510 AtomicCommitOutcome::Ingested { 511 projection: Some(_), 512 .. 513 } 514 )); 515 516 let replay = store.commit(request.clone()).await.expect("replay"); 517 assert_eq!(replay.disposition(), AtomicCommitDisposition::Replay); 518 assert_eq!(replay.outcome(), committed.outcome()); 519 let reconstructed = store 520 .receipt(request.commit_id()) 521 .await 522 .expect("receipt lookup") 523 .expect("receipt"); 524 assert_eq!(reconstructed.outcome(), committed.outcome()); 525 526 let conflict = commit( 527 1, 528 2, 529 AtomicWorkflow::Ingested(Box::new(CommitIngested::new(admission(90), None))), 530 ); 531 assert_eq!( 532 store.commit(conflict).await, 533 Err(Error::AtomicCommitConflict) 534 ); 535 } 536 537 #[tokio::test] 538 async fn read_only_mode_rejects_atomic_mutation() { 539 let store = store(EventStoreMode::ReadOnly).await; 540 let request = commit( 541 2, 542 2, 543 AtomicWorkflow::Ingested(Box::new(CommitIngested::new(admission(90), None))), 544 ); 545 assert_eq!(store.commit(request).await, Err(Error::BackendUnavailable)); 546 } 547 548 #[test] 549 fn receipt_decoder_rejects_retired_and_trailing_payloads() { 550 assert_eq!( 551 decode_outcome(&vec![0; RECEIPT_MAX_BYTES + 1]), 552 Err(Error::AtomicCommitFailed) 553 ); 554 assert_eq!(decode_outcome(&[0]), Err(Error::AtomicCommitFailed)); 555 assert_eq!( 556 decode_outcome(&[RECEIPT_FORMAT_VERSION, 0]), 557 Err(Error::AtomicCommitFailed) 558 ); 559 let outcome = AtomicCommitOutcome::Ingested { 560 admission: AdmissionReceipt::new( 561 radroots_storage::event::EventId::from_bytes([1; 32]), 562 EventPosition::new( 563 SourceGeneration::new([2; 32]).expect("generation"), 564 EventSequence::new(1).expect("sequence"), 565 ), 566 AdmissionStage::Raw, 567 AdmissionDisposition::Inserted, 568 ), 569 projection: None, 570 }; 571 let mut bytes = encode_outcome(&outcome).expect("encode"); 572 bytes.push(0); 573 assert_eq!(decode_outcome(&bytes), Err(Error::AtomicCommitFailed)); 574 } 575 }