authored_draft_submission_tests.rs (22129B)
1 use super::*; 2 use crate::{OpenMode, OpenOptions, Paths}; 3 use radroots_storage::{ 4 authored_draft::{ 5 AuthoredDraft, AuthoredDraftId, AuthoredDraftRevision, AuthoredDraftStage, 6 AuthoredDraftStore, 7 }, 8 authored_draft_query::AuthoredDraftScope, 9 authored_draft_submission::{AuthoredDraftSource, PrepareFromDraft}, 10 event::SourceGeneration, 11 memory::MemoryStorage, 12 }; 13 use sha2::{Digest, Sha256}; 14 use tempfile::TempDir; 15 16 async fn open(temp: &TempDir, mode: OpenMode) -> SqliteStorage { 17 let options = OpenOptions::new(Paths::from_directory(temp.path()).unwrap(), mode); 18 let options = if mode == OpenMode::Create { 19 options 20 .with_source_generation(SourceGeneration::new([9; 32]).unwrap(), 9) 21 .unwrap() 22 } else { 23 options 24 }; 25 SqliteStorage::open(options).await.unwrap() 26 } 27 fn fixture() -> (AuthoredDraft, PrepareFromDraft) { 28 fixture_stage(AuthoredDraftStage::Queued) 29 } 30 fn fixture_stage(stage: AuthoredDraftStage) -> (AuthoredDraft, PrepareFromDraft) { 31 let (ordinary, plan) = super::tests::prepare(); 32 let AuthoredAtomicCommand::Prepare(preparation) = ordinary else { 33 unreachable!() 34 }; 35 let scope = AuthoredDraftScope::new([3; 32]).unwrap(); 36 let source = AuthoredDraft::initial( 37 AuthoredDraftId::new([8; 16]).unwrap(), 38 *plan.author().as_bytes(), 39 "fixture.partial.v1", 40 b"partial 0.".to_vec(), 41 AuthoredDraftStage::Draft, 42 None, 43 9, 44 ) 45 .unwrap() 46 .with_scope(scope) 47 .unwrap(); 48 let payload = b"complete captured request".to_vec(); 49 let intent = AuthoredDraft::reconstruct( 50 AuthoredDraftId::new([7; 16]).unwrap(), 51 AuthoredDraftRevision::INITIAL, 52 *source.author(), 53 "fixture.intent.v1", 54 payload.clone(), 55 Sha256::digest(&payload).into(), 56 stage, 57 matches!( 58 stage, 59 AuthoredDraftStage::ReadyToSign | AuthoredDraftStage::Queued 60 ) 61 .then_some(preparation.operation().operation_id()), 62 10, 63 10, 64 ) 65 .unwrap() 66 .with_scope(scope) 67 .unwrap(); 68 let constructor = if intent.operation_id().is_some() { 69 PrepareFromDraft::new 70 } else { 71 PrepareFromDraft::new_waiting 72 }; 73 let request = constructor( 74 AtomicCommitId::new([4; 16]).unwrap(), 75 AuthoredDraftSource::capture(&source).unwrap(), 76 intent, 77 preparation, 78 ) 79 .unwrap(); 80 (source, request) 81 } 82 fn command(request: &PrepareFromDraft) -> AuthoredAtomicCommand { 83 AuthoredAtomicCommand::PrepareFromDraft(Box::new(request.clone())) 84 } 85 async fn empty_submission( 86 store: &SqliteStorage, 87 source: &AuthoredDraft, 88 request: &PrepareFromDraft, 89 ) { 90 for table in [ 91 "operations", 92 "artifacts", 93 "delivery_plans", 94 "delivery_targets", 95 "atomic_commits", 96 ] { 97 // Audited test-only identifiers come exclusively from the literal list above. 98 let count: i64 = sqlx::query_scalar(sqlx::AssertSqlSafe(format!( 99 "SELECT COUNT(*) FROM radroots_runtime_authored_{table}" 100 ))) 101 .fetch_one(store.pool()) 102 .await 103 .unwrap(); 104 assert_eq!(count, 0, "{table}"); 105 } 106 assert_eq!( 107 store 108 .authored_draft_heads(*source.author(), 10) 109 .await 110 .unwrap(), 111 std::slice::from_ref(source) 112 ); 113 assert!( 114 store 115 .authored_receipt(request.commit_id()) 116 .await 117 .unwrap() 118 .is_none() 119 ); 120 } 121 122 #[tokio::test] 123 async fn submission_survives_lost_callback_reopen_and_matches_memory_after_later_save() { 124 for stage in [ 125 AuthoredDraftStage::Queued, 126 AuthoredDraftStage::ReadyToSign, 127 AuthoredDraftStage::Draft, 128 AuthoredDraftStage::MediaPreparing, 129 AuthoredDraftStage::MediaUploading, 130 ] { 131 replay_after_progress(stage).await; 132 } 133 } 134 async fn replay_after_progress(stage: AuthoredDraftStage) { 135 let temp = TempDir::new().unwrap(); 136 let store = open(&temp, OpenMode::Create).await; 137 let memory = MemoryStorage::default(); 138 let (source, request) = fixture_stage(stage); 139 store 140 .append_authored_draft(source.clone(), None) 141 .await 142 .unwrap(); 143 memory 144 .append_authored_draft(source.clone(), None) 145 .await 146 .unwrap(); 147 let expected = memory.execute_authored(command(&request)).await.unwrap(); 148 assert_eq!( 149 store.execute_authored(command(&request)).await.unwrap(), 150 expected 151 ); 152 if request.intent().operation_id().is_none() { 153 let progress = request 154 .intent() 155 .successor( 156 b"verified prerequisite status".to_vec(), 157 AuthoredDraftStage::ReadyToSign, 158 Some(request.preparation().operation().operation_id()), 159 11, 160 ) 161 .unwrap(); 162 for target in [&store as &dyn AuthoredDraftStore, &memory] { 163 target 164 .append_authored_draft(progress.clone(), Some(request.intent().revision())) 165 .await 166 .unwrap(); 167 } 168 } 169 let newer = source 170 .successor(b"later edit".to_vec(), AuthoredDraftStage::Draft, None, 11) 171 .unwrap(); 172 for target in [&store as &dyn AuthoredDraftStore, &memory] { 173 target 174 .append_authored_draft(newer.clone(), Some(source.revision())) 175 .await 176 .unwrap(); 177 } 178 // Discard the callback value and release the actual writer before recovery. 179 store.close().await.unwrap(); 180 assert!(store.execute_authored(command(&request)).await.is_err()); 181 let store = open(&temp, OpenMode::ReadWriteExisting).await; 182 let replay = store.execute_authored(command(&request)).await.unwrap(); 183 assert_eq!( 184 replay, 185 memory.execute_authored(command(&request)).await.unwrap() 186 ); 187 assert_eq!(replay.disposition(), AtomicCommitDisposition::Replay); 188 assert_eq!(replay.outcome(), expected.outcome()); 189 assert_eq!( 190 store 191 .authored_draft_head(request.intent().draft_id()) 192 .await 193 .unwrap(), 194 memory 195 .authored_draft_head(request.intent().draft_id()) 196 .await 197 .unwrap() 198 ); 199 let ordinary = AuthoredAtomicCommand::Prepare(request.preparation().clone()); 200 assert_eq!( 201 store.execute_authored(ordinary.clone()).await.unwrap(), 202 memory.execute_authored(ordinary).await.unwrap() 203 ); 204 // A valid changed context under the same author/command conflicts before CAS. 205 let mut changed = serde_json::to_value(&request).unwrap(); 206 changed["source"]["scope"] = serde_json::json!([6; 32].to_vec()); 207 changed["intent"]["scope"] = serde_json::json!([6; 32].to_vec()); 208 let changed = serde_json::from_value::<PrepareFromDraft>(changed).unwrap(); 209 assert_eq!(command(&changed).digest(), command(&request).digest()); 210 assert_eq!( 211 store.execute_authored(command(&changed)).await, 212 Err(Error::AtomicCommitConflict) 213 ); 214 let reader = open(&temp, OpenMode::ReadOnly).await; 215 assert!( 216 reader 217 .authored_receipt(request.commit_id()) 218 .await 219 .unwrap() 220 .is_some() 221 ); 222 assert!(reader.execute_authored(command(&request)).await.is_err()); 223 reader.close().await.unwrap(); 224 store.close().await.unwrap(); 225 } 226 227 #[tokio::test] 228 async fn submission_rolls_back_before_and_after_every_record_and_receipt_insert() { 229 for stage in [ 230 AuthoredDraftStage::Queued, 231 AuthoredDraftStage::Draft, 232 AuthoredDraftStage::MediaPreparing, 233 AuthoredDraftStage::MediaUploading, 234 ] { 235 rollback_at_every_write(stage).await; 236 } 237 } 238 async fn rollback_at_every_write(stage: AuthoredDraftStage) { 239 let (source, request) = fixture_stage(stage); 240 let ordinary = AuthoredAtomicCommand::Prepare(request.preparation().clone()); 241 let hex = |id: AtomicCommitId| { 242 id.as_bytes() 243 .iter() 244 .map(|b| format!("{b:02x}")) 245 .collect::<String>() 246 }; 247 let points = [ 248 ("operations", String::new()), 249 ("artifacts", String::new()), 250 ("delivery_plans", String::new()), 251 ("delivery_targets", "WHEN NEW.ordinal = 0".into()), 252 ("delivery_targets", "WHEN NEW.ordinal = 1".into()), 253 ( 254 "draft_revisions", 255 "WHEN NEW.draft_id != x'08080808080808080808080808080808'".into(), 256 ), 257 ( 258 "atomic_commits", 259 format!("WHEN NEW.commit_id = x'{}'", hex(ordinary.commit_id())), 260 ), 261 ( 262 "atomic_commits", 263 format!("WHEN NEW.commit_id = x'{}'", hex(request.commit_id())), 264 ), 265 ]; 266 for timing in ["BEFORE", "AFTER"] { 267 for (table, condition) in &points { 268 let temp = TempDir::new().unwrap(); 269 let store = open(&temp, OpenMode::Create).await; 270 store 271 .append_authored_draft(source.clone(), None) 272 .await 273 .unwrap(); 274 // Timing/table/condition are fixed test literals plus hex-encoded typed IDs. 275 sqlx::query(sqlx::AssertSqlSafe(format!("CREATE TRIGGER submission_fault {timing} INSERT ON radroots_runtime_authored_{table} {condition} BEGIN SELECT RAISE(ABORT, 'injected submission failure'); END"))) 276 .execute(store.pool()).await.unwrap(); 277 assert!( 278 store.execute_authored(command(&request)).await.is_err(), 279 "{timing} {table} {condition}" 280 ); 281 empty_submission(&store, &source, &request).await; 282 sqlx::query("DROP TRIGGER submission_fault") 283 .execute(store.pool()) 284 .await 285 .unwrap(); 286 store.close().await.unwrap(); 287 let store = open(&temp, OpenMode::ReadWriteExisting).await; 288 empty_submission(&store, &source, &request).await; 289 assert_eq!( 290 store 291 .execute_authored(command(&request)) 292 .await 293 .unwrap() 294 .disposition(), 295 AtomicCommitDisposition::Committed 296 ); 297 store.close().await.unwrap(); 298 } 299 } 300 } 301 302 #[tokio::test] 303 async fn abandoned_precommit_and_sqlite_full_preserve_source_without_success_association() { 304 for stage in [ 305 AuthoredDraftStage::Queued, 306 AuthoredDraftStage::MediaPreparing, 307 ] { 308 abandoned_and_full(stage).await; 309 } 310 } 311 async fn abandoned_and_full(stage: AuthoredDraftStage) { 312 let temp = TempDir::new().unwrap(); 313 let store = open(&temp, OpenMode::Create).await; 314 let (source, request) = fixture_stage(stage); 315 store 316 .append_authored_draft(source.clone(), None) 317 .await 318 .unwrap(); 319 let mut transaction = store.pool().begin_with("BEGIN IMMEDIATE").await.unwrap(); 320 let receipt = execute_transaction(&mut transaction, &command(&request)) 321 .await 322 .unwrap(); 323 assert_eq!(receipt.disposition(), AtomicCommitDisposition::Committed); 324 // This internal result has not crossed the public commit boundary. 325 drop(transaction); 326 empty_submission(&store, &source, &request).await; 327 let mut transaction = store.pool().begin_with("BEGIN IMMEDIATE").await.unwrap(); 328 let page_count: i64 = sqlx::query_scalar("PRAGMA page_count") 329 .fetch_one(&mut *transaction) 330 .await 331 .unwrap(); 332 // The PRAGMA accepts an integer, not a bind parameter; this is a typed i64. 333 sqlx::query(sqlx::AssertSqlSafe(format!( 334 "PRAGMA max_page_count = {page_count}" 335 ))) 336 .execute(&mut *transaction) 337 .await 338 .unwrap(); 339 let mut large = serde_json::to_value(&request).unwrap(); 340 let payload = vec![97u8; 100_000]; 341 large["intent"]["payload"] = serde_json::json!(payload); 342 large["intent"]["payload_sha256"] = serde_json::json!(Sha256::digest(&payload).to_vec()); 343 let large: PrepareFromDraft = serde_json::from_value(large).unwrap(); 344 assert_eq!( 345 execute_transaction(&mut transaction, &command(&large)).await, 346 Err(Error::SpaceInsufficient) 347 ); 348 // SQLITE_FULL may already roll back the transaction at the engine boundary. 349 let _ = transaction.rollback().await; 350 empty_submission(&store, &source, &request).await; 351 store.close().await.unwrap(); 352 let store = open(&temp, OpenMode::ReadWriteExisting).await; 353 assert_eq!( 354 store 355 .execute_authored(command(&request)) 356 .await 357 .unwrap() 358 .disposition(), 359 AtomicCommitDisposition::Committed 360 ); 361 store.close().await.unwrap(); 362 } 363 364 #[tokio::test] 365 async fn concurrent_duplicates_and_save_have_one_atomic_order() { 366 for stage in [ 367 AuthoredDraftStage::Queued, 368 AuthoredDraftStage::MediaPreparing, 369 ] { 370 concurrent_submit_and_save(stage).await; 371 } 372 } 373 async fn concurrent_submit_and_save(stage: AuthoredDraftStage) { 374 for _ in 0..4 { 375 let temp = TempDir::new().unwrap(); 376 let store = open(&temp, OpenMode::Create).await; 377 let (source, request) = fixture_stage(stage); 378 store 379 .append_authored_draft(source.clone(), None) 380 .await 381 .unwrap(); 382 let newer = source 383 .successor( 384 b"concurrent edit".to_vec(), 385 AuthoredDraftStage::Draft, 386 None, 387 11, 388 ) 389 .unwrap(); 390 let (first, second, save) = tokio::join!( 391 store.execute_authored(command(&request)), 392 store.execute_authored(command(&request)), 393 store.append_authored_draft(newer.clone(), Some(source.revision())) 394 ); 395 save.unwrap(); 396 match (first, second) { 397 (Ok(a), Ok(b)) => { 398 assert_eq!(a.outcome(), b.outcome()); 399 assert_ne!(a.disposition(), b.disposition()); 400 assert_eq!( 401 store 402 .execute_authored(command(&request)) 403 .await 404 .unwrap() 405 .disposition(), 406 AtomicCommitDisposition::Replay 407 ); 408 let operations: i64 = 409 sqlx::query_scalar("SELECT COUNT(*) FROM radroots_runtime_authored_operations") 410 .fetch_one(store.pool()) 411 .await 412 .unwrap(); 413 assert_eq!(operations, 1); 414 } 415 (Err(Error::DraftRevisionConflict), Err(Error::DraftRevisionConflict)) => { 416 assert!( 417 store 418 .authored_receipt(request.commit_id()) 419 .await 420 .unwrap() 421 .is_none() 422 ); 423 assert!( 424 store 425 .authored_draft_head(request.intent().draft_id()) 426 .await 427 .unwrap() 428 .is_none() 429 ); 430 } 431 other => panic!("non-atomic submit/save result: {other:?}"), 432 } 433 assert_eq!( 434 store.authored_draft_head(source.draft_id()).await.unwrap(), 435 Some(newer) 436 ); 437 store.close().await.unwrap(); 438 } 439 } 440 441 #[tokio::test] 442 async fn commit_failure_returns_no_success_and_rolls_back_every_authored_write() { 443 for stage in [ 444 AuthoredDraftStage::Queued, 445 AuthoredDraftStage::MediaPreparing, 446 ] { 447 commit_failure(stage).await; 448 } 449 } 450 async fn commit_failure(stage: AuthoredDraftStage) { 451 let temp = TempDir::new().unwrap(); 452 let store = open(&temp, OpenMode::Create).await; 453 let (source, request) = fixture_stage(stage); 454 store 455 .append_authored_draft(source.clone(), None) 456 .await 457 .unwrap(); 458 sqlx::query("CREATE TABLE submission_commit_fault (parent BLOB REFERENCES radroots_runtime_authored_operations(operation_id) DEFERRABLE INITIALLY DEFERRED)") 459 .execute(store.pool()).await.unwrap(); 460 // The last receipt succeeds; only the actual transaction COMMIT fails. 461 sqlx::query("CREATE TRIGGER submission_commit_fault_trigger AFTER INSERT ON radroots_runtime_authored_atomic_commits BEGIN INSERT INTO submission_commit_fault VALUES (x'99999999999999999999999999999999'); END") 462 .execute(store.pool()).await.unwrap(); 463 assert!(store.execute_authored(command(&request)).await.is_err()); 464 empty_submission(&store, &source, &request).await; 465 sqlx::query("DROP TRIGGER submission_commit_fault_trigger") 466 .execute(store.pool()) 467 .await 468 .unwrap(); 469 sqlx::query("DROP TABLE submission_commit_fault") 470 .execute(store.pool()) 471 .await 472 .unwrap(); 473 assert_eq!( 474 store 475 .execute_authored(command(&request)) 476 .await 477 .unwrap() 478 .disposition(), 479 AtomicCommitDisposition::Committed 480 ); 481 store.close().await.unwrap(); 482 } 483 484 #[test] 485 fn storage_submission_machine_policy_matches_native_bounds() { 486 let policy: serde_json::Value = serde_json::from_str(include_str!( 487 "../../../contracts/architecture/decisions/authored_draft_submission.v1.json" 488 )) 489 .unwrap(); 490 let bounds = &policy["bounds"]; 491 assert_eq!( 492 bounds["page_records"], 493 radroots_storage::authored_draft::AUTHORED_DRAFT_QUERY_LIMIT_MAX 494 ); 495 assert_eq!( 496 bounds["decoded_page_payload_bytes"], 497 radroots_storage::authored_draft_query::AUTHORED_DRAFT_PAGE_PAYLOAD_MAX_BYTES 498 ); 499 assert_eq!( 500 bounds["serialized_page_snapshot_bytes"], 501 radroots_storage::authored_draft_query::AUTHORED_DRAFT_PAGE_SNAPSHOT_MAX_BYTES 502 ); 503 assert_eq!( 504 bounds["existing_atomic_receipt_snapshot_bytes"], 505 super::SNAPSHOT_MAX_BYTES 506 ); 507 assert_eq!( 508 policy["owners"], 509 serde_json::json!([ 510 "radroots_storage", 511 "radroots_storage_sqlite", 512 "radroots_sync" 513 ]) 514 ); 515 } 516 517 #[tokio::test] 518 async fn existing_intent_and_partial_graph_id_collisions_cannot_create_associations() { 519 let (source, request) = fixture(); 520 let temp = TempDir::new().unwrap(); 521 let store = open(&temp, OpenMode::Create).await; 522 store 523 .append_authored_draft(source.clone(), None) 524 .await 525 .unwrap(); 526 store 527 .append_authored_draft(request.intent().clone(), None) 528 .await 529 .unwrap(); 530 assert_eq!( 531 store.execute_authored(command(&request)).await, 532 Err(Error::DraftRevisionConflict) 533 ); 534 assert!( 535 store 536 .authored_receipt(request.commit_id()) 537 .await 538 .unwrap() 539 .is_none() 540 ); 541 assert!( 542 store 543 .authored_operation(request.preparation().operation().operation_id()) 544 .await 545 .unwrap() 546 .is_none() 547 ); 548 store.close().await.unwrap(); 549 let temp = TempDir::new().unwrap(); 550 let store = open(&temp, OpenMode::Create).await; 551 store 552 .execute_authored(AuthoredAtomicCommand::Prepare( 553 request.preparation().clone(), 554 )) 555 .await 556 .unwrap(); 557 for artifact_collision in [true, false] { 558 let base = request.preparation(); 559 let operation_id = OperationInstanceId::new([20; 16]).unwrap(); 560 let artifact_id = if artifact_collision { 561 base.artifacts()[0].artifact_id() 562 } else { 563 AuthoredArtifactId::new([21; 16]).unwrap() 564 }; 565 let plan_id = if artifact_collision { 566 AuthoredDeliveryPlanId::new([22; 16]).unwrap() 567 } else { 568 base.delivery_plans()[0].plan_id() 569 }; 570 let artifact = AuthoredArtifact::planned( 571 artifact_id, 572 operation_id, 573 0, 574 base.artifacts()[0].plan().unwrap().decode().unwrap().plan(), 575 10, 576 ) 577 .unwrap(); 578 let preparation = radroots_storage::authored_atomic::PrepareAuthoredOperation::new( 579 AuthoredOperation::new(operation_id, vec![artifact_id], 10).unwrap(), 580 vec![artifact], 581 vec![ 582 AuthoredDeliveryPlan::new( 583 plan_id, 584 artifact_id, 585 base.delivery_plans()[0].intent().clone(), 586 10, 587 ) 588 .unwrap(), 589 ], 590 base.input_digest(), 591 10, 592 ) 593 .unwrap(); 594 assert_eq!( 595 store 596 .execute_authored(AuthoredAtomicCommand::Prepare(preparation)) 597 .await, 598 Err(Error::AtomicCommitConflict) 599 ); 600 assert!( 601 store 602 .authored_operation(operation_id) 603 .await 604 .unwrap() 605 .is_none() 606 ); 607 } 608 store.close().await.unwrap(); 609 } 610 611 #[tokio::test] 612 async fn submission_receipt_shadow_times_must_match_the_captured_request() { 613 let temp = TempDir::new().unwrap(); 614 let store = open(&temp, OpenMode::Create).await; 615 let (source, request) = fixture(); 616 store.append_authored_draft(source, None).await.unwrap(); 617 store.execute_authored(command(&request)).await.unwrap(); 618 for requested in [9i64, 11] { 619 let row = sqlx::query("SELECT commit_id, commit_digest, ? AS requested_at_unix_ms, committed_at_unix_ms, receipt FROM radroots_runtime_authored_atomic_commits WHERE commit_id = ?") 620 .bind(requested).bind(request.commit_id().as_bytes().as_slice()).fetch_one(store.pool()).await.unwrap(); 621 assert_eq!(decode_receipt_row(&row), Err(Error::AtomicCommitFailed)); 622 } 623 assert_eq!( 624 store 625 .execute_authored(command(&request)) 626 .await 627 .unwrap() 628 .disposition(), 629 AtomicCommitDisposition::Replay 630 ); 631 store.close().await.unwrap(); 632 }