authored.rs (102056B)
1 use crate::SqliteStorage; 2 use crate::backend::map_backend; 3 #[path = "authored_delivery_facts.rs"] 4 mod delivery_facts; 5 #[path = "authored_delivery_reconciliation.rs"] 6 mod delivery_reconciliation; 7 use radroots_storage::{ 8 Error, 9 atomic::{AtomicCommitDigest, AtomicCommitDisposition, AtomicCommitId}, 10 authored::{ 11 AdmissionState, ArtifactOrigin, AuthoredArtifact, AuthoredArtifactId, AuthoredOperation, 12 FailureClass, SigningState, WorkFailure, WorkPhase, 13 }, 14 authored_atomic::{ 15 AuthoredAtomicCommand, AuthoredAtomicOutcome, AuthoredAtomicReceipt, AuthoredAtomicStorage, 16 AuthoredWorkTarget, CancelAuthoredTarget, ClaimAuthoredTarget, WorkFence, 17 }, 18 authored_delivery::{ 19 AuthoredDeliveryPlan, AuthoredDeliveryPlanId, AuthoredDeliveryState, DeliveryAttemptOutcome, 20 }, 21 event::BoxFuture, 22 journal::OperationInstanceId, 23 }; 24 use radroots_transport::{SinkFailure, outcome::Retryability, policy::SatisfactionState}; 25 use serde::{Deserialize, Serialize, de::DeserializeOwned}; 26 use sqlx::{Row, Sqlite, sqlite::SqliteRow}; 27 28 const SNAPSHOT_MAX_BYTES: usize = 4 * 1024 * 1024; 29 30 #[derive(Serialize, Deserialize)] 31 struct ReceiptSnapshot { 32 outcome: AuthoredAtomicOutcome, 33 } 34 35 impl AuthoredAtomicStorage for SqliteStorage { 36 fn authored_delivery_history( 37 &self, 38 plan_id: AuthoredDeliveryPlanId, 39 ) -> BoxFuture< 40 '_, 41 Result<Option<radroots_storage::authored_delivery::AuthoredDeliveryHistory>, Error>, 42 > { 43 Box::pin(async move { 44 let mut transaction = self.pool().begin().await.map_err(map_backend)?; 45 let result = delivery_reconciliation::history(&mut transaction, plan_id).await; 46 let rollback = transaction.rollback().await.map_err(map_backend); 47 match result { 48 Ok(history) => { 49 rollback?; 50 Ok(history) 51 } 52 Err(error) => Err(error), 53 } 54 }) 55 } 56 fn execute_authored( 57 &self, 58 command: AuthoredAtomicCommand, 59 ) -> BoxFuture<'_, Result<AuthoredAtomicReceipt, Error>> { 60 Box::pin(async move { 61 self.require_authored_writer()?; 62 let mut transaction = self 63 .pool() 64 .begin_with("BEGIN IMMEDIATE") 65 .await 66 .map_err(map_backend)?; 67 let result = execute_transaction(&mut transaction, &command).await; 68 match result { 69 Ok(receipt) => { 70 transaction.commit().await.map_err(map_backend)?; 71 Ok(receipt) 72 } 73 Err(primary) => { 74 let _ = transaction.rollback().await; 75 Err(primary) 76 } 77 } 78 }) 79 } 80 81 fn authored_receipt( 82 &self, 83 commit_id: AtomicCommitId, 84 ) -> BoxFuture<'_, Result<Option<AuthoredAtomicReceipt>, Error>> { 85 Box::pin(async move { 86 sqlx::query( 87 "SELECT commit_id, commit_digest, requested_at_unix_ms, 88 committed_at_unix_ms, receipt 89 FROM radroots_runtime_authored_atomic_commits WHERE commit_id = ?", 90 ) 91 .bind(commit_id.as_bytes().as_slice()) 92 .fetch_optional(self.pool()) 93 .await 94 .map_err(map_backend)? 95 .as_ref() 96 .map(decode_receipt_row) 97 .transpose() 98 }) 99 } 100 101 fn authored_operation( 102 &self, 103 operation_id: OperationInstanceId, 104 ) -> BoxFuture<'_, Result<Option<AuthoredOperation>, Error>> { 105 Box::pin(async move { 106 let Some(row) = sqlx::query( 107 "SELECT operation_id, artifact_count, created_at_unix_ms, 108 updated_at_unix_ms, revision, snapshot 109 FROM radroots_runtime_authored_operations WHERE operation_id = ?", 110 ) 111 .bind(operation_id.as_bytes().as_slice()) 112 .fetch_optional(self.pool()) 113 .await 114 .map_err(map_backend)? 115 else { 116 return Ok(None); 117 }; 118 let operation = decode_operation_row(&row)?; 119 for (ordinal, artifact_id) in operation.artifact_ids().iter().enumerate() { 120 let artifact = sqlx::query( 121 "SELECT * FROM radroots_runtime_authored_artifacts WHERE artifact_id = ?", 122 ) 123 .bind(artifact_id.as_bytes().as_slice()) 124 .fetch_optional(self.pool()) 125 .await 126 .map_err(map_backend)? 127 .as_ref() 128 .map(decode_artifact_row) 129 .transpose()? 130 .ok_or(Error::InvalidAuthoredOperation)?; 131 if artifact.operation_id() != operation.operation_id() 132 || usize::from(artifact.ordinal()) != ordinal 133 { 134 return Err(Error::InvalidAuthoredOperation); 135 } 136 } 137 Ok(Some(operation)) 138 }) 139 } 140 141 fn authored_artifact( 142 &self, 143 artifact_id: AuthoredArtifactId, 144 ) -> BoxFuture<'_, Result<Option<AuthoredArtifact>, Error>> { 145 Box::pin(async move { 146 let Some(row) = sqlx::query( 147 "SELECT * FROM radroots_runtime_authored_artifacts WHERE artifact_id = ?", 148 ) 149 .bind(artifact_id.as_bytes().as_slice()) 150 .fetch_optional(self.pool()) 151 .await 152 .map_err(map_backend)? 153 else { 154 return Ok(None); 155 }; 156 let artifact = decode_artifact_row(&row)?; 157 let operation_row = sqlx::query( 158 "SELECT operation_id, artifact_count, created_at_unix_ms, 159 updated_at_unix_ms, revision, snapshot 160 FROM radroots_runtime_authored_operations WHERE operation_id = ?", 161 ) 162 .bind(artifact.operation_id().as_bytes().as_slice()) 163 .fetch_optional(self.pool()) 164 .await 165 .map_err(map_backend)? 166 .ok_or(Error::InvalidAuthoredArtifact)?; 167 let operation = decode_operation_row(&operation_row)?; 168 if operation 169 .artifact_ids() 170 .get(usize::from(artifact.ordinal())) 171 != Some(&artifact.artifact_id()) 172 { 173 return Err(Error::InvalidAuthoredArtifact); 174 } 175 Ok(Some(artifact)) 176 }) 177 } 178 179 fn authored_delivery_plan( 180 &self, 181 plan_id: AuthoredDeliveryPlanId, 182 ) -> BoxFuture<'_, Result<Option<AuthoredDeliveryPlan>, Error>> { 183 Box::pin(async move { load_plan_pool(self, plan_id).await }) 184 } 185 } 186 187 impl SqliteStorage { 188 fn require_authored_writer(&self) -> Result<(), Error> { 189 if self.event_mode() == radroots_storage::status::EventStoreMode::ReadOnly { 190 return Err(Error::BackendUnavailable); 191 } 192 Ok(()) 193 } 194 } 195 196 async fn execute_transaction( 197 transaction: &mut sqlx::Transaction<'_, Sqlite>, 198 command: &AuthoredAtomicCommand, 199 ) -> Result<AuthoredAtomicReceipt, Error> { 200 if let Some(row) = sqlx::query( 201 "SELECT commit_id, commit_digest, requested_at_unix_ms, 202 committed_at_unix_ms, receipt 203 FROM radroots_runtime_authored_atomic_commits WHERE commit_id = ?", 204 ) 205 .bind(command.commit_id().as_bytes().as_slice()) 206 .fetch_optional(&mut **transaction) 207 .await 208 .map_err(map_backend)? 209 { 210 let committed = decode_receipt_row(&row)?; 211 if !committed.matches_command(command) { 212 return Err(Error::AtomicCommitConflict); 213 } 214 return AuthoredAtomicReceipt::from_durable_parts( 215 committed.commit_id(), 216 committed.digest(), 217 AtomicCommitDisposition::Replay, 218 committed.committed_at_unix_ms(), 219 committed.outcome().clone(), 220 ); 221 } 222 223 let outcome = execute_command(transaction, command.clone()).await?; 224 commit_outcome(transaction, command, outcome).await 225 } 226 227 async fn commit_outcome( 228 transaction: &mut sqlx::Transaction<'_, Sqlite>, 229 command: &AuthoredAtomicCommand, 230 outcome: AuthoredAtomicOutcome, 231 ) -> Result<AuthoredAtomicReceipt, Error> { 232 let committed_at = match (command, &outcome) { 233 (AuthoredAtomicCommand::RecordDelivery(_), AuthoredAtomicOutcome::DeliveryPlan(plan)) => { 234 command 235 .requested_at_unix_ms() 236 .max(plan.updated_at_unix_ms()) 237 } 238 (AuthoredAtomicCommand::RecordSigned(_), AuthoredAtomicOutcome::Artifact(artifact)) => { 239 command 240 .requested_at_unix_ms() 241 .max(artifact.updated_at_unix_ms()) 242 } 243 _ => command.requested_at_unix_ms(), 244 }; 245 let receipt = AuthoredAtomicReceipt::new( 246 command, 247 AtomicCommitDisposition::Committed, 248 committed_at, 249 outcome, 250 )?; 251 sqlx::query( 252 "INSERT INTO radroots_runtime_authored_atomic_commits ( 253 commit_id, commit_digest, phase, target_id, requested_at_unix_ms, 254 committed_at_unix_ms, receipt 255 ) VALUES (?, ?, ?, ?, ?, ?, ?)", 256 ) 257 .bind(receipt.commit_id().as_bytes().as_slice()) 258 .bind(receipt.digest().as_bytes().as_slice()) 259 .bind(command_phase(command)) 260 .bind(command_target(command).as_slice()) 261 .bind(i64_from_u64(command.requested_at_unix_ms())?) 262 .bind(i64_from_u64(receipt.committed_at_unix_ms())?) 263 .bind(encode_snapshot(&ReceiptSnapshot { 264 outcome: receipt.outcome().clone(), 265 })?) 266 .execute(&mut **transaction) 267 .await 268 .map_err(map_backend)?; 269 delivery_reconciliation::record_claim(transaction, command).await?; 270 Ok(receipt) 271 } 272 273 async fn prepare_operation( 274 transaction: &mut sqlx::Transaction<'_, Sqlite>, 275 value: radroots_storage::authored_atomic::PrepareAuthoredOperation, 276 ) -> Result<AuthoredAtomicOutcome, Error> { 277 if row_exists( 278 transaction, 279 "SELECT 1 FROM radroots_runtime_authored_operations WHERE operation_id = ?", 280 value.operation().operation_id().as_bytes(), 281 ) 282 .await? 283 || any_artifact_exists(transaction, value.artifacts()).await? 284 || any_plan_exists(transaction, value.delivery_plans()).await? 285 { 286 return Err(Error::AtomicCommitConflict); 287 } 288 persist_operation(transaction, value.operation()).await?; 289 for artifact in value.artifacts() { 290 persist_artifact(transaction, artifact).await?; 291 } 292 for plan in value.delivery_plans() { 293 persist_plan(transaction, plan).await?; 294 } 295 Ok(AuthoredAtomicOutcome::Prepared { 296 operation: value.operation().clone(), 297 artifacts: value.artifacts().to_vec(), 298 delivery_plans: value.delivery_plans().to_vec(), 299 }) 300 } 301 302 async fn execute_command( 303 transaction: &mut sqlx::Transaction<'_, Sqlite>, 304 command: AuthoredAtomicCommand, 305 ) -> Result<AuthoredAtomicOutcome, Error> { 306 match command { 307 AuthoredAtomicCommand::Prepare(value) => prepare_operation(transaction, value).await, 308 AuthoredAtomicCommand::PrepareFromDraft(value) => { 309 value.validate()?; 310 let head = 311 crate::authored_draft::load_head_tx(transaction, value.source().draft_id()).await?; 312 if !head 313 .as_ref() 314 .is_some_and(|head| value.source().matches(head)) 315 { 316 return Err(Error::DraftRevisionConflict); 317 } 318 if crate::authored_draft::load_head_tx(transaction, value.intent().draft_id()) 319 .await? 320 .is_some() 321 { 322 return Err(Error::DraftRevisionConflict); 323 } 324 let ordinary = AuthoredAtomicCommand::Prepare(value.preparation().clone()); 325 let prepared = prepare_operation(transaction, value.preparation().clone()).await?; 326 crate::authored_draft::insert_draft_tx(transaction, value.intent()).await?; 327 commit_outcome(transaction, &ordinary, prepared).await?; 328 Ok(AuthoredAtomicOutcome::Submitted(value)) 329 } 330 AuthoredAtomicCommand::Claim(value) => match value.target() { 331 ClaimAuthoredTarget::ArtifactSigning(id) => { 332 let mut artifact = load_artifact_tx(transaction, *id).await?; 333 artifact.set_signing_claim( 334 value.claim().clone(), 335 value.claim().acquired_at_unix_ms(), 336 )?; 337 persist_artifact(transaction, &artifact).await?; 338 Ok(AuthoredAtomicOutcome::Artifact(artifact)) 339 } 340 ClaimAuthoredTarget::ArtifactAdmission(id) => { 341 let mut artifact = load_artifact_tx(transaction, *id).await?; 342 artifact.set_admission_claim( 343 value.claim().clone(), 344 value.claim().acquired_at_unix_ms(), 345 )?; 346 persist_artifact(transaction, &artifact).await?; 347 Ok(AuthoredAtomicOutcome::Artifact(artifact)) 348 } 349 ClaimAuthoredTarget::DeliveryPlan(id) => { 350 let mut plan = load_plan_tx(transaction, *id).await?; 351 let artifact = load_artifact_tx(transaction, plan.artifact_id()).await?; 352 if artifact.signing_state() != SigningState::Signed { 353 return Err(Error::InvalidAuthoredTransition); 354 } 355 plan.claim(value.claim().clone(), value.claim().acquired_at_unix_ms())?; 356 persist_plan(transaction, &plan).await?; 357 Ok(AuthoredAtomicOutcome::DeliveryPlan(plan)) 358 } 359 }, 360 AuthoredAtomicCommand::RecordSigned(value) => { 361 let row = sqlx::query( 362 "SELECT commit_id, commit_digest, requested_at_unix_ms, 363 committed_at_unix_ms, receipt 364 FROM radroots_runtime_authored_atomic_commits WHERE commit_id = ?", 365 ) 366 .bind(value.claim_command().commit_id().as_bytes().as_slice()) 367 .fetch_optional(&mut **transaction) 368 .await 369 .map_err(map_backend)? 370 .ok_or(Error::AtomicWorkflowMismatch)?; 371 let original = decode_receipt_row(&row)?; 372 let mut artifact = load_artifact_tx(transaction, value.artifact_id()).await?; 373 let already_signed = artifact.signed().is_some(); 374 value.apply_to(&mut artifact, &original)?; 375 persist_artifact(transaction, &artifact).await?; 376 if !already_signed && artifact.signing_state() == SigningState::Signed { 377 let plan_ids = sqlx::query_scalar::<_, Vec<u8>>( 378 "SELECT plan_id FROM radroots_runtime_authored_delivery_plans 379 WHERE artifact_id = ? ORDER BY plan_id", 380 ) 381 .bind(value.artifact_id().as_bytes().as_slice()) 382 .fetch_all(&mut **transaction) 383 .await 384 .map_err(map_backend)?; 385 for bytes in plan_ids { 386 let id = AuthoredDeliveryPlanId::new(array(bytes)?)?; 387 let mut plan = load_plan_tx(transaction, id).await?; 388 if !plan.state().is_terminal() { 389 plan.bind_signed_event( 390 value.event().clone(), 391 value.observed_at_unix_ms().max(plan.updated_at_unix_ms()), 392 )?; 393 persist_plan(transaction, &plan).await?; 394 } 395 } 396 } 397 Ok(AuthoredAtomicOutcome::Artifact(artifact)) 398 } 399 AuthoredAtomicCommand::ApplySigned(value) => { 400 let mut artifact = load_artifact_tx(transaction, value.artifact_id()).await?; 401 require_artifact_claim( 402 artifact.signing_claim(), 403 value.fence(), 404 value.applied_at_unix_ms(), 405 )?; 406 artifact.record_signed(value.event().clone(), value.applied_at_unix_ms())?; 407 persist_artifact(transaction, &artifact).await?; 408 let plan_ids = sqlx::query_scalar::<_, Vec<u8>>( 409 "SELECT plan_id FROM radroots_runtime_authored_delivery_plans 410 WHERE artifact_id = ? ORDER BY plan_id", 411 ) 412 .bind(value.artifact_id().as_bytes().as_slice()) 413 .fetch_all(&mut **transaction) 414 .await 415 .map_err(map_backend)?; 416 for bytes in plan_ids { 417 let id = AuthoredDeliveryPlanId::new(array(bytes)?)?; 418 let mut plan = load_plan_tx(transaction, id).await?; 419 plan.bind_signed_event(value.event().clone(), value.applied_at_unix_ms())?; 420 persist_plan(transaction, &plan).await?; 421 } 422 Ok(AuthoredAtomicOutcome::Artifact(artifact)) 423 } 424 AuthoredAtomicCommand::ApplyAdmission(value) => { 425 let mut artifact = load_artifact_tx(transaction, value.artifact_id()).await?; 426 require_artifact_claim( 427 artifact.admission_claim(), 428 value.fence(), 429 value.applied_at_unix_ms(), 430 )?; 431 artifact.record_admission( 432 value.state(), 433 value.failure().cloned(), 434 value.retry().cloned(), 435 value.applied_at_unix_ms(), 436 )?; 437 persist_artifact(transaction, &artifact).await?; 438 Ok(AuthoredAtomicOutcome::Artifact(artifact)) 439 } 440 AuthoredAtomicCommand::RecordDelivery(value) => { 441 let row = sqlx::query( 442 "SELECT commit_id, commit_digest, requested_at_unix_ms, 443 committed_at_unix_ms, receipt 444 FROM radroots_runtime_authored_atomic_commits WHERE commit_id = ?", 445 ) 446 .bind(value.claim_command().commit_id().as_bytes().as_slice()) 447 .fetch_optional(&mut **transaction) 448 .await 449 .map_err(map_backend)? 450 .ok_or(Error::AtomicWorkflowMismatch)?; 451 let original = decode_receipt_row(&row)?; 452 let mut plan = load_plan_tx(transaction, value.plan_id()).await?; 453 value.apply_to(&mut plan, &original)?; 454 persist_plan(transaction, &plan).await?; 455 Ok(AuthoredAtomicOutcome::DeliveryPlan(plan)) 456 } 457 AuthoredAtomicCommand::ReconcileDelivery(value) => { 458 let plan = delivery_reconciliation::reconcile(transaction, &value).await?; 459 Ok(AuthoredAtomicOutcome::DeliveryPlan(plan)) 460 } 461 AuthoredAtomicCommand::ApplyDelivery(value) => { 462 let mut plan = load_plan_tx(transaction, value.plan_id()).await?; 463 match value.outcome().clone() { 464 DeliveryAttemptOutcome::Receipt(receipt) => plan.apply_receipt( 465 value.fence().token(), 466 value.fence().generation(), 467 value.fence().row_revision(), 468 receipt, 469 value.retry().cloned(), 470 value.applied_at_unix_ms(), 471 )?, 472 DeliveryAttemptOutcome::SinkFailure(failure) => plan.apply_sink_failure( 473 value.fence().token(), 474 value.fence().generation(), 475 value.fence().row_revision(), 476 failure, 477 value.retry().cloned(), 478 value.applied_at_unix_ms(), 479 )?, 480 } 481 persist_plan(transaction, &plan).await?; 482 Ok(AuthoredAtomicOutcome::DeliveryPlan(plan)) 483 } 484 AuthoredAtomicCommand::ApplyFailure(value) => match value.target() { 485 AuthoredWorkTarget::Artifact(id) => { 486 let mut artifact = load_artifact_tx(transaction, *id).await?; 487 apply_artifact_failure(&mut artifact, &value)?; 488 persist_artifact(transaction, &artifact).await?; 489 Ok(AuthoredAtomicOutcome::Artifact(artifact)) 490 } 491 AuthoredWorkTarget::DeliveryPlan(id) => { 492 let mut plan = load_plan_tx(transaction, *id).await?; 493 if value.failure().phase() != WorkPhase::Delivery 494 || value.failure().class() == FailureClass::Indeterminate 495 { 496 return Err(Error::AtomicWorkflowMismatch); 497 } 498 let retryability = match value.failure().class() { 499 FailureClass::Retryable => Retryability::Retryable, 500 FailureClass::Terminal => Retryability::Terminal, 501 FailureClass::Indeterminate => unreachable!(), 502 }; 503 let failure = SinkFailure::for_request( 504 plan.request().ok_or(Error::InvalidAuthoredDeliveryPlan)?, 505 value.failure().code(), 506 retryability, 507 value.failure().retry_after_unix_ms(), 508 value.failure().diagnostic().map(str::to_owned), 509 Vec::new(), 510 ) 511 .map_err(|_| Error::AtomicWorkflowMismatch)?; 512 plan.apply_sink_failure( 513 value.fence().token(), 514 value.fence().generation(), 515 value.fence().row_revision(), 516 failure, 517 value.retry().cloned(), 518 value.applied_at_unix_ms(), 519 )?; 520 persist_plan(transaction, &plan).await?; 521 Ok(AuthoredAtomicOutcome::DeliveryPlan(plan)) 522 } 523 }, 524 AuthoredAtomicCommand::Cancel(value) => match value.target() { 525 CancelAuthoredTarget::ArtifactSigning(id) => { 526 let mut artifact = load_artifact_tx(transaction, *id).await?; 527 require_revision(artifact.revision().get(), value.expected_revision().get())?; 528 artifact.cancel_signing(value.cancelled_at_unix_ms())?; 529 persist_artifact(transaction, &artifact).await?; 530 Ok(AuthoredAtomicOutcome::Artifact(artifact)) 531 } 532 CancelAuthoredTarget::ArtifactAdmission(id) => { 533 let mut artifact = load_artifact_tx(transaction, *id).await?; 534 require_revision(artifact.revision().get(), value.expected_revision().get())?; 535 artifact.record_admission( 536 AdmissionState::Cancelled, 537 Some(WorkFailure::new( 538 "cancelled", 539 WorkPhase::Admission, 540 FailureClass::Terminal, 541 None, 542 None, 543 )?), 544 None, 545 value.cancelled_at_unix_ms(), 546 )?; 547 persist_artifact(transaction, &artifact).await?; 548 Ok(AuthoredAtomicOutcome::Artifact(artifact)) 549 } 550 CancelAuthoredTarget::DeliveryPlan(id) => { 551 let mut plan = load_plan_tx(transaction, *id).await?; 552 require_revision(plan.revision().get(), value.expected_revision().get())?; 553 plan.request_stop(value.cancelled_at_unix_ms())?; 554 persist_plan(transaction, &plan).await?; 555 Ok(AuthoredAtomicOutcome::DeliveryPlan(plan)) 556 } 557 }, 558 } 559 } 560 561 fn apply_artifact_failure( 562 artifact: &mut AuthoredArtifact, 563 value: &radroots_storage::authored_atomic::ApplyWorkFailure, 564 ) -> Result<(), Error> { 565 match value.failure().phase() { 566 WorkPhase::Signing => { 567 require_artifact_claim( 568 artifact.signing_claim(), 569 value.fence(), 570 value.applied_at_unix_ms(), 571 )?; 572 artifact.record_signing_failure( 573 value.failure().clone(), 574 value.retry().cloned(), 575 value.applied_at_unix_ms(), 576 ) 577 } 578 WorkPhase::Admission => { 579 require_artifact_claim( 580 artifact.admission_claim(), 581 value.fence(), 582 value.applied_at_unix_ms(), 583 )?; 584 let state = match value.failure().class() { 585 FailureClass::Retryable => AdmissionState::Retryable, 586 FailureClass::Terminal => AdmissionState::Rejected, 587 FailureClass::Indeterminate => return Err(Error::InvalidAuthoredTransition), 588 }; 589 artifact.record_admission( 590 state, 591 Some(value.failure().clone()), 592 value.retry().cloned(), 593 value.applied_at_unix_ms(), 594 ) 595 } 596 WorkPhase::Delivery => Err(Error::AtomicWorkflowMismatch), 597 } 598 } 599 600 pub(crate) async fn persist_operation( 601 transaction: &mut sqlx::Transaction<'_, Sqlite>, 602 operation: &AuthoredOperation, 603 ) -> Result<(), Error> { 604 sqlx::query( 605 "INSERT INTO radroots_runtime_authored_operations ( 606 operation_id, artifact_count, created_at_unix_ms, updated_at_unix_ms, 607 revision, snapshot 608 ) VALUES (?, ?, ?, ?, ?, ?)", 609 ) 610 .bind(operation.operation_id().as_bytes().as_slice()) 611 .bind(i64::try_from(operation.artifact_ids().len()).map_err(|_| Error::AtomicCommitFailed)?) 612 .bind(i64_from_u64(operation.created_at_unix_ms())?) 613 .bind(i64_from_u64(operation.updated_at_unix_ms())?) 614 .bind(i64_from_u64(operation.revision().get())?) 615 .bind(encode_snapshot(operation)?) 616 .execute(&mut **transaction) 617 .await 618 .map_err(map_backend)?; 619 Ok(()) 620 } 621 622 pub(crate) async fn persist_artifact( 623 transaction: &mut sqlx::Transaction<'_, Sqlite>, 624 artifact: &AuthoredArtifact, 625 ) -> Result<(), Error> { 626 let signing = artifact.signing_claim(); 627 let admission = artifact.admission_claim(); 628 let retry_not_before = artifact 629 .signing_retry() 630 .or_else(|| artifact.admission_retry()) 631 .map(|retry| retry.not_before_unix_ms()); 632 sqlx::query( 633 "INSERT INTO radroots_runtime_authored_artifacts ( 634 artifact_id, operation_id, ordinal, origin, signing_state, admission_state, 635 plan_wire, signed_raw_json, signed_raw_sha256, 636 signing_claim_token, signing_claim_generation, signing_claim_revision, 637 signing_claim_expires_at_unix_ms, admission_claim_token, 638 admission_claim_generation, admission_claim_revision, 639 admission_claim_expires_at_unix_ms, retry_not_before_unix_ms, 640 last_failure_code, created_at_unix_ms, updated_at_unix_ms, revision, snapshot, 641 signing_stop 642 ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) 643 ON CONFLICT(artifact_id) DO UPDATE SET 644 operation_id=excluded.operation_id, ordinal=excluded.ordinal, origin=excluded.origin, 645 signing_state=excluded.signing_state, admission_state=excluded.admission_state, 646 plan_wire=excluded.plan_wire, signed_raw_json=excluded.signed_raw_json, 647 signed_raw_sha256=excluded.signed_raw_sha256, 648 signing_claim_token=excluded.signing_claim_token, 649 signing_claim_generation=excluded.signing_claim_generation, 650 signing_claim_revision=excluded.signing_claim_revision, 651 signing_claim_expires_at_unix_ms=excluded.signing_claim_expires_at_unix_ms, 652 admission_claim_token=excluded.admission_claim_token, 653 admission_claim_generation=excluded.admission_claim_generation, 654 admission_claim_revision=excluded.admission_claim_revision, 655 admission_claim_expires_at_unix_ms=excluded.admission_claim_expires_at_unix_ms, 656 retry_not_before_unix_ms=excluded.retry_not_before_unix_ms, 657 last_failure_code=excluded.last_failure_code, 658 updated_at_unix_ms=excluded.updated_at_unix_ms, revision=excluded.revision, 659 snapshot=excluded.snapshot, signing_stop=excluded.signing_stop", 660 ) 661 .bind(artifact.artifact_id().as_bytes().as_slice()) 662 .bind(artifact.operation_id().as_bytes().as_slice()) 663 .bind(i64::from(artifact.ordinal())) 664 .bind(origin_name(artifact.origin())) 665 .bind(physical_signing_name(artifact)) 666 .bind(admission_name(artifact.admission_state())) 667 .bind(artifact.plan().map(|plan| plan.wire_json())) 668 .bind( 669 artifact 670 .signed() 671 .map(|signed| signed.event().raw_json().as_bytes()), 672 ) 673 .bind( 674 artifact 675 .signed() 676 .map(|signed| signed.raw_json_sha256().as_slice()), 677 ) 678 .bind(signing.map(|claim| claim.token().as_slice())) 679 .bind( 680 signing 681 .map(|claim| i64_from_u64(claim.generation().get())) 682 .transpose()?, 683 ) 684 .bind( 685 signing 686 .map(|claim| i64_from_u64(claim.row_revision().get())) 687 .transpose()?, 688 ) 689 .bind( 690 signing 691 .map(|claim| i64_from_u64(claim.expires_at_unix_ms())) 692 .transpose()?, 693 ) 694 .bind(admission.map(|claim| claim.token().as_slice())) 695 .bind( 696 admission 697 .map(|claim| i64_from_u64(claim.generation().get())) 698 .transpose()?, 699 ) 700 .bind( 701 admission 702 .map(|claim| i64_from_u64(claim.row_revision().get())) 703 .transpose()?, 704 ) 705 .bind( 706 admission 707 .map(|claim| i64_from_u64(claim.expires_at_unix_ms())) 708 .transpose()?, 709 ) 710 .bind(retry_not_before.map(i64_from_u64).transpose()?) 711 .bind(artifact.last_failure().map(WorkFailure::code)) 712 .bind(i64_from_u64(artifact.created_at_unix_ms())?) 713 .bind(i64_from_u64(artifact.updated_at_unix_ms())?) 714 .bind(i64_from_u64(artifact.revision().get())?) 715 .bind(encode_snapshot(artifact)?) 716 .bind(signing_stop(artifact)) 717 .execute(&mut **transaction) 718 .await 719 .map_err(map_backend)?; 720 Ok(()) 721 } 722 723 pub(crate) async fn persist_plan( 724 transaction: &mut sqlx::Transaction<'_, Sqlite>, 725 plan: &AuthoredDeliveryPlan, 726 ) -> Result<(), Error> { 727 persist_plan_v11(transaction, plan).await?; 728 delivery_facts::persist(transaction, plan).await?; 729 delivery_reconciliation::persist(transaction, plan).await 730 } 731 732 // Legacy conversion runs before the delivery-facts forward migration. 733 pub(crate) async fn persist_plan_v11( 734 transaction: &mut sqlx::Transaction<'_, Sqlite>, 735 plan: &AuthoredDeliveryPlan, 736 ) -> Result<(), Error> { 737 let claim = plan.claim_evidence(); 738 sqlx::query( 739 "INSERT INTO radroots_runtime_authored_delivery_plans ( 740 plan_id, artifact_id, request_digest, state, attempt_count, 741 claim_token, claim_generation, claim_revision, claim_expires_at_unix_ms, 742 retry_not_before_unix_ms, last_failure_code, created_at_unix_ms, 743 updated_at_unix_ms, revision, snapshot 744 ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) 745 ON CONFLICT(plan_id) DO UPDATE SET 746 artifact_id=excluded.artifact_id, request_digest=excluded.request_digest, 747 state=excluded.state, attempt_count=excluded.attempt_count, 748 claim_token=excluded.claim_token, claim_generation=excluded.claim_generation, 749 claim_revision=excluded.claim_revision, 750 claim_expires_at_unix_ms=excluded.claim_expires_at_unix_ms, 751 retry_not_before_unix_ms=excluded.retry_not_before_unix_ms, 752 last_failure_code=excluded.last_failure_code, 753 updated_at_unix_ms=excluded.updated_at_unix_ms, revision=excluded.revision, 754 snapshot=excluded.snapshot", 755 ) 756 .bind(plan.plan_id().as_bytes().as_slice()) 757 .bind(plan.artifact_id().as_bytes().as_slice()) 758 .bind(plan.request_digest().as_slice()) 759 .bind(delivery_state_name(plan.state())) 760 .bind(i64::from(plan.attempt_count())) 761 .bind(claim.map(|value| value.token().as_slice())) 762 .bind( 763 claim 764 .map(|value| i64_from_u64(value.generation().get())) 765 .transpose()?, 766 ) 767 .bind( 768 claim 769 .map(|value| i64_from_u64(value.row_revision().get())) 770 .transpose()?, 771 ) 772 .bind( 773 claim 774 .map(|value| i64_from_u64(value.expires_at_unix_ms())) 775 .transpose()?, 776 ) 777 .bind( 778 plan.retry() 779 .map(|retry| i64_from_u64(retry.not_before_unix_ms())) 780 .transpose()?, 781 ) 782 .bind(plan.last_failure().map(WorkFailure::code)) 783 .bind(i64_from_u64(plan.created_at_unix_ms())?) 784 .bind(i64_from_u64(plan.updated_at_unix_ms())?) 785 .bind(i64_from_u64(plan.revision().get())?) 786 .bind(encode_snapshot(plan)?) 787 .execute(&mut **transaction) 788 .await 789 .map_err(map_backend)?; 790 791 sqlx::query("DELETE FROM radroots_runtime_authored_delivery_targets WHERE plan_id = ?") 792 .bind(plan.plan_id().as_bytes().as_slice()) 793 .execute(&mut **transaction) 794 .await 795 .map_err(map_backend)?; 796 sqlx::query("DELETE FROM radroots_runtime_authored_delivery_attempts WHERE plan_id = ?") 797 .bind(plan.plan_id().as_bytes().as_slice()) 798 .execute(&mut **transaction) 799 .await 800 .map_err(map_backend)?; 801 for (ordinal, target) in plan.intent().target_set().targets().iter().enumerate() { 802 sqlx::query( 803 "INSERT INTO radroots_runtime_authored_delivery_targets ( 804 plan_id, ordinal, target_fingerprint, target_snapshot 805 ) VALUES (?, ?, ?, ?)", 806 ) 807 .bind(plan.plan_id().as_bytes().as_slice()) 808 .bind(i64::try_from(ordinal).map_err(|_| Error::InvalidAuthoredDeliveryPlan)?) 809 .bind(target.fingerprint().as_str()) 810 .bind(encode_snapshot(target)?) 811 .execute(&mut **transaction) 812 .await 813 .map_err(map_backend)?; 814 } 815 for attempt in plan.attempts() { 816 sqlx::query( 817 "INSERT INTO radroots_runtime_authored_delivery_attempts ( 818 plan_id, attempt, satisfaction, recorded_at_unix_ms, outcome_snapshot 819 ) VALUES (?, ?, ?, ?, ?)", 820 ) 821 .bind(plan.plan_id().as_bytes().as_slice()) 822 .bind(i64::from(attempt.attempt().get())) 823 .bind(satisfaction_name(attempt.satisfaction())) 824 .bind(i64_from_u64(attempt.recorded_at_unix_ms())?) 825 .bind(encode_snapshot(attempt.outcome())?) 826 .execute(&mut **transaction) 827 .await 828 .map_err(map_backend)?; 829 } 830 Ok(()) 831 } 832 833 async fn load_artifact_tx( 834 transaction: &mut sqlx::Transaction<'_, Sqlite>, 835 artifact_id: AuthoredArtifactId, 836 ) -> Result<AuthoredArtifact, Error> { 837 sqlx::query("SELECT * FROM radroots_runtime_authored_artifacts WHERE artifact_id = ?") 838 .bind(artifact_id.as_bytes().as_slice()) 839 .fetch_optional(&mut **transaction) 840 .await 841 .map_err(map_backend)? 842 .as_ref() 843 .map(decode_artifact_row) 844 .transpose()? 845 .ok_or(Error::InvalidAuthoredArtifact) 846 } 847 848 async fn load_plan_tx( 849 transaction: &mut sqlx::Transaction<'_, Sqlite>, 850 plan_id: AuthoredDeliveryPlanId, 851 ) -> Result<AuthoredDeliveryPlan, Error> { 852 load_optional_plan_tx(transaction, plan_id) 853 .await? 854 .ok_or(Error::InvalidAuthoredDeliveryPlan) 855 } 856 857 async fn load_optional_plan_tx( 858 transaction: &mut sqlx::Transaction<'_, Sqlite>, 859 plan_id: AuthoredDeliveryPlanId, 860 ) -> Result<Option<AuthoredDeliveryPlan>, Error> { 861 let Some(row) = 862 sqlx::query("SELECT * FROM radroots_runtime_authored_delivery_plans WHERE plan_id = ?") 863 .bind(plan_id.as_bytes().as_slice()) 864 .fetch_optional(&mut **transaction) 865 .await 866 .map_err(map_backend)? 867 else { 868 return Ok(None); 869 }; 870 let plan = decode_plan_row(&row)?; 871 validate_plan_children_tx(transaction, &plan).await?; 872 Ok(Some(plan)) 873 } 874 875 async fn load_plan_pool( 876 storage: &SqliteStorage, 877 plan_id: AuthoredDeliveryPlanId, 878 ) -> Result<Option<AuthoredDeliveryPlan>, Error> { 879 // Plan and append-only facts must come from the same read snapshot even 880 // when a late result commits between the individual SELECT statements. 881 let mut transaction = storage.pool().begin().await.map_err(map_backend)?; 882 let result = load_optional_plan_tx(&mut transaction, plan_id).await; 883 let rollback = transaction.rollback().await.map_err(map_backend); 884 match result { 885 Ok(plan) => { 886 rollback?; 887 Ok(plan) 888 } 889 Err(error) => Err(error), 890 } 891 } 892 893 async fn validate_plan_children_tx( 894 transaction: &mut sqlx::Transaction<'_, Sqlite>, 895 plan: &AuthoredDeliveryPlan, 896 ) -> Result<(), Error> { 897 let targets = sqlx::query( 898 "SELECT ordinal, target_fingerprint, target_snapshot 899 FROM radroots_runtime_authored_delivery_targets WHERE plan_id = ? ORDER BY ordinal", 900 ) 901 .bind(plan.plan_id().as_bytes().as_slice()) 902 .fetch_all(&mut **transaction) 903 .await 904 .map_err(map_backend)?; 905 let attempts = sqlx::query( 906 "SELECT attempt, satisfaction, recorded_at_unix_ms, outcome_snapshot 907 FROM radroots_runtime_authored_delivery_attempts WHERE plan_id = ? ORDER BY attempt", 908 ) 909 .bind(plan.plan_id().as_bytes().as_slice()) 910 .fetch_all(&mut **transaction) 911 .await 912 .map_err(map_backend)?; 913 validate_plan_children(plan, &targets, &attempts)?; 914 let facts = sqlx::query(delivery_facts::SELECT) 915 .bind(plan.plan_id().as_bytes().as_slice()) 916 .fetch_all(&mut **transaction) 917 .await 918 .map_err(map_backend)?; 919 delivery_facts::validate(plan, &facts)?; 920 delivery_reconciliation::validate(transaction, plan).await 921 } 922 923 fn validate_plan_children( 924 plan: &AuthoredDeliveryPlan, 925 targets: &[SqliteRow], 926 attempts: &[SqliteRow], 927 ) -> Result<(), Error> { 928 if targets.len() != plan.intent().target_set().targets().len() 929 || attempts.len() != plan.attempts().len() 930 { 931 return Err(Error::InvalidAuthoredDeliveryPlan); 932 } 933 for (ordinal, (row, expected)) in targets 934 .iter() 935 .zip(plan.intent().target_set().targets()) 936 .enumerate() 937 { 938 let decoded = 939 decode_snapshot::<radroots_transport::Target>(column(row, "target_snapshot")?)?; 940 if column::<i64>(row, "ordinal")? != i64::try_from(ordinal).unwrap_or(i64::MAX) 941 || column::<String>(row, "target_fingerprint")? != expected.fingerprint().as_str() 942 || decoded != *expected 943 { 944 return Err(Error::InvalidAuthoredDeliveryPlan); 945 } 946 } 947 for (row, expected) in attempts.iter().zip(plan.attempts()) { 948 let decoded = decode_snapshot::<DeliveryAttemptOutcome>(column(row, "outcome_snapshot")?)?; 949 if column::<i64>(row, "attempt")? != i64::from(expected.attempt().get()) 950 || column::<String>(row, "satisfaction")? != satisfaction_name(expected.satisfaction()) 951 || u64_from_i64(column(row, "recorded_at_unix_ms")?)? != expected.recorded_at_unix_ms() 952 || decoded != *expected.outcome() 953 { 954 return Err(Error::InvalidAuthoredDeliveryPlan); 955 } 956 } 957 Ok(()) 958 } 959 960 fn decode_operation_row(row: &SqliteRow) -> Result<AuthoredOperation, Error> { 961 let value = decode_snapshot::<AuthoredOperation>(column(row, "snapshot")?)?; 962 if column::<Vec<u8>>(row, "operation_id")?.as_slice() != value.operation_id().as_bytes() 963 || column::<i64>(row, "artifact_count")? 964 != i64::try_from(value.artifact_ids().len()) 965 .map_err(|_| Error::InvalidAuthoredOperation)? 966 || u64_from_i64(column(row, "created_at_unix_ms")?)? != value.created_at_unix_ms() 967 || u64_from_i64(column(row, "updated_at_unix_ms")?)? != value.updated_at_unix_ms() 968 || u64_from_i64(column(row, "revision")?)? != value.revision().get() 969 { 970 return Err(Error::InvalidAuthoredOperation); 971 } 972 Ok(value) 973 } 974 975 fn decode_artifact_row(row: &SqliteRow) -> Result<AuthoredArtifact, Error> { 976 let value = decode_snapshot::<AuthoredArtifact>(column(row, "snapshot")?)?; 977 let plan_wire = column::<Option<Vec<u8>>>(row, "plan_wire")?; 978 let raw = column::<Option<Vec<u8>>>(row, "signed_raw_json")?; 979 let raw_digest = column::<Option<Vec<u8>>>(row, "signed_raw_sha256")?; 980 let signing_claim = value.signing_claim(); 981 let admission_claim = value.admission_claim(); 982 let retry_not_before = value 983 .signing_retry() 984 .or_else(|| value.admission_retry()) 985 .map(|retry| retry.not_before_unix_ms()); 986 if column::<Vec<u8>>(row, "artifact_id")?.as_slice() != value.artifact_id().as_bytes() 987 || column::<Vec<u8>>(row, "operation_id")?.as_slice() != value.operation_id().as_bytes() 988 || column::<i64>(row, "ordinal")? != i64::from(value.ordinal()) 989 || column::<String>(row, "origin")? != origin_name(value.origin()) 990 || column::<String>(row, "signing_state")? != physical_signing_name(&value) 991 || column::<Option<String>>(row, "signing_stop")?.as_deref() != signing_stop(&value) 992 || column::<String>(row, "admission_state")? != admission_name(value.admission_state()) 993 || plan_wire.as_deref() != value.plan().map(|plan| plan.wire_json()) 994 || raw.as_deref() 995 != value 996 .signed() 997 .map(|signed| signed.event().raw_json().as_bytes()) 998 || raw_digest.as_deref() 999 != value 1000 .signed() 1001 .map(|signed| signed.raw_json_sha256().as_slice()) 1002 || !claim_columns_match(row, "signing", signing_claim)? 1003 || !claim_columns_match(row, "admission", admission_claim)? 1004 || column::<Option<i64>>(row, "retry_not_before_unix_ms")? 1005 != retry_not_before.map(i64_from_u64).transpose()? 1006 || column::<Option<String>>(row, "last_failure_code")?.as_deref() 1007 != value.last_failure().map(WorkFailure::code) 1008 || u64_from_i64(column(row, "created_at_unix_ms")?)? != value.created_at_unix_ms() 1009 || u64_from_i64(column(row, "updated_at_unix_ms")?)? != value.updated_at_unix_ms() 1010 || u64_from_i64(column(row, "revision")?)? != value.revision().get() 1011 { 1012 return Err(Error::InvalidAuthoredArtifact); 1013 } 1014 Ok(value) 1015 } 1016 1017 fn decode_plan_row(row: &SqliteRow) -> Result<AuthoredDeliveryPlan, Error> { 1018 let value = decode_snapshot::<AuthoredDeliveryPlan>(column(row, "snapshot")?)?; 1019 if column::<Vec<u8>>(row, "plan_id")?.as_slice() != value.plan_id().as_bytes() 1020 || column::<Vec<u8>>(row, "artifact_id")?.as_slice() != value.artifact_id().as_bytes() 1021 || column::<Vec<u8>>(row, "request_digest")?.as_slice() != value.request_digest() 1022 || column::<String>(row, "state")? != delivery_state_name(value.state()) 1023 || column::<i64>(row, "attempt_count")? != i64::from(value.attempt_count()) 1024 || column::<Option<i64>>(row, "stop_requested_at_unix_ms")? 1025 != value 1026 .stop_requested_at_unix_ms() 1027 .map(i64_from_u64) 1028 .transpose()? 1029 || !claim_columns_match(row, "claim", value.claim_evidence())? 1030 || column::<Option<i64>>(row, "retry_not_before_unix_ms")? 1031 != value 1032 .retry() 1033 .map(|retry| i64_from_u64(retry.not_before_unix_ms())) 1034 .transpose()? 1035 || column::<Option<String>>(row, "last_failure_code")?.as_deref() 1036 != value.last_failure().map(WorkFailure::code) 1037 || u64_from_i64(column(row, "created_at_unix_ms")?)? != value.created_at_unix_ms() 1038 || u64_from_i64(column(row, "updated_at_unix_ms")?)? != value.updated_at_unix_ms() 1039 || u64_from_i64(column(row, "revision")?)? != value.revision().get() 1040 { 1041 return Err(Error::InvalidAuthoredDeliveryPlan); 1042 } 1043 Ok(value) 1044 } 1045 1046 fn claim_columns_match( 1047 row: &SqliteRow, 1048 prefix: &str, 1049 claim: Option<&radroots_storage::authored::WorkClaim>, 1050 ) -> Result<bool, Error> { 1051 let (token_column, generation_column, revision_column, expiry_column) = match prefix { 1052 "signing" => ( 1053 "signing_claim_token", 1054 "signing_claim_generation", 1055 "signing_claim_revision", 1056 "signing_claim_expires_at_unix_ms", 1057 ), 1058 "admission" => ( 1059 "admission_claim_token", 1060 "admission_claim_generation", 1061 "admission_claim_revision", 1062 "admission_claim_expires_at_unix_ms", 1063 ), 1064 "claim" => ( 1065 "claim_token", 1066 "claim_generation", 1067 "claim_revision", 1068 "claim_expires_at_unix_ms", 1069 ), 1070 _ => return Err(Error::AtomicCommitFailed), 1071 }; 1072 Ok(column::<Option<Vec<u8>>>(row, token_column)?.as_deref() 1073 == claim.map(|value| value.token().as_slice()) 1074 && column::<Option<i64>>(row, generation_column)? 1075 == claim 1076 .map(|value| i64_from_u64(value.generation().get())) 1077 .transpose()? 1078 && column::<Option<i64>>(row, revision_column)? 1079 == claim 1080 .map(|value| i64_from_u64(value.row_revision().get())) 1081 .transpose()? 1082 && column::<Option<i64>>(row, expiry_column)? 1083 == claim 1084 .map(|value| i64_from_u64(value.expires_at_unix_ms())) 1085 .transpose()?) 1086 } 1087 1088 fn decode_receipt_row(row: &SqliteRow) -> Result<AuthoredAtomicReceipt, Error> { 1089 let commit_id = AtomicCommitId::new(array(column(row, "commit_id")?)?)?; 1090 let digest = AtomicCommitDigest::new(array(column(row, "commit_digest")?)?); 1091 let requested = u64_from_i64(column(row, "requested_at_unix_ms")?)?; 1092 let committed = u64_from_i64(column(row, "committed_at_unix_ms")?)?; 1093 let snapshot = decode_snapshot::<ReceiptSnapshot>(column(row, "receipt")?)?; 1094 if committed < requested 1095 || matches!(&snapshot.outcome, AuthoredAtomicOutcome::Submitted(value) if value.preparation().requested_at_unix_ms() != requested) 1096 { 1097 return Err(Error::AtomicCommitFailed); 1098 } 1099 AuthoredAtomicReceipt::from_durable_parts( 1100 commit_id, 1101 digest, 1102 AtomicCommitDisposition::Committed, 1103 committed, 1104 snapshot.outcome, 1105 ) 1106 .map_err(|_| Error::AtomicCommitFailed) 1107 } 1108 1109 async fn any_artifact_exists( 1110 transaction: &mut sqlx::Transaction<'_, Sqlite>, 1111 artifacts: &[AuthoredArtifact], 1112 ) -> Result<bool, Error> { 1113 for artifact in artifacts { 1114 if row_exists( 1115 transaction, 1116 "SELECT 1 FROM radroots_runtime_authored_artifacts WHERE artifact_id = ?", 1117 artifact.artifact_id().as_bytes(), 1118 ) 1119 .await? 1120 { 1121 return Ok(true); 1122 } 1123 } 1124 Ok(false) 1125 } 1126 1127 async fn any_plan_exists( 1128 transaction: &mut sqlx::Transaction<'_, Sqlite>, 1129 plans: &[AuthoredDeliveryPlan], 1130 ) -> Result<bool, Error> { 1131 for plan in plans { 1132 if row_exists( 1133 transaction, 1134 "SELECT 1 FROM radroots_runtime_authored_delivery_plans WHERE plan_id = ?", 1135 plan.plan_id().as_bytes(), 1136 ) 1137 .await? 1138 { 1139 return Ok(true); 1140 } 1141 } 1142 Ok(false) 1143 } 1144 1145 async fn row_exists<const N: usize>( 1146 transaction: &mut sqlx::Transaction<'_, Sqlite>, 1147 query: &'static str, 1148 id: &[u8; N], 1149 ) -> Result<bool, Error> { 1150 Ok(sqlx::query_scalar::<_, i64>(query) 1151 .bind(id.as_slice()) 1152 .fetch_optional(&mut **transaction) 1153 .await 1154 .map_err(map_backend)? 1155 .is_some()) 1156 } 1157 1158 fn require_artifact_claim( 1159 claim: Option<&radroots_storage::authored::WorkClaim>, 1160 fence: &WorkFence, 1161 now_unix_ms: u64, 1162 ) -> Result<(), Error> { 1163 if !claim.is_some_and(|claim| { 1164 claim.matches_fence( 1165 fence.token(), 1166 fence.generation(), 1167 fence.row_revision(), 1168 now_unix_ms, 1169 ) 1170 }) { 1171 return Err(Error::DeliveryPlanClaimConflict); 1172 } 1173 Ok(()) 1174 } 1175 1176 fn require_revision(actual: u64, expected: u64) -> Result<(), Error> { 1177 if actual != expected { 1178 return Err(Error::InvalidAuthoredTransition); 1179 } 1180 Ok(()) 1181 } 1182 1183 pub(crate) fn encode_snapshot<T: Serialize>(value: &T) -> Result<Vec<u8>, Error> { 1184 let bytes = serde_json::to_vec(value).map_err(|_| Error::AtomicCommitFailed)?; 1185 if bytes.len() < 2 || bytes.len() > SNAPSHOT_MAX_BYTES { 1186 return Err(Error::AtomicCommitFailed); 1187 } 1188 Ok(bytes) 1189 } 1190 1191 fn decode_snapshot<T: DeserializeOwned>(bytes: Vec<u8>) -> Result<T, Error> { 1192 if bytes.len() < 2 || bytes.len() > SNAPSHOT_MAX_BYTES { 1193 return Err(Error::AtomicCommitFailed); 1194 } 1195 serde_json::from_slice(&bytes).map_err(|_| Error::AtomicCommitFailed) 1196 } 1197 1198 fn command_target(command: &AuthoredAtomicCommand) -> [u8; 16] { 1199 match command { 1200 AuthoredAtomicCommand::PrepareFromDraft(value) => { 1201 *value.preparation().operation().operation_id().as_bytes() 1202 } 1203 AuthoredAtomicCommand::Prepare(value) => *value.operation().operation_id().as_bytes(), 1204 AuthoredAtomicCommand::Claim(value) => match value.target() { 1205 ClaimAuthoredTarget::ArtifactSigning(id) 1206 | ClaimAuthoredTarget::ArtifactAdmission(id) => *id.as_bytes(), 1207 ClaimAuthoredTarget::DeliveryPlan(id) => *id.as_bytes(), 1208 }, 1209 AuthoredAtomicCommand::ApplySigned(value) => *value.artifact_id().as_bytes(), 1210 AuthoredAtomicCommand::RecordSigned(value) => *value.artifact_id().as_bytes(), 1211 AuthoredAtomicCommand::RecordDelivery(value) => *value.plan_id().as_bytes(), 1212 AuthoredAtomicCommand::ReconcileDelivery(value) => *value.plan_id().as_bytes(), 1213 AuthoredAtomicCommand::ApplyAdmission(value) => *value.artifact_id().as_bytes(), 1214 AuthoredAtomicCommand::ApplyDelivery(value) => *value.plan_id().as_bytes(), 1215 AuthoredAtomicCommand::ApplyFailure(value) => match value.target() { 1216 AuthoredWorkTarget::Artifact(id) => *id.as_bytes(), 1217 AuthoredWorkTarget::DeliveryPlan(id) => *id.as_bytes(), 1218 }, 1219 AuthoredAtomicCommand::Cancel(value) => match value.target() { 1220 CancelAuthoredTarget::ArtifactSigning(id) 1221 | CancelAuthoredTarget::ArtifactAdmission(id) => *id.as_bytes(), 1222 CancelAuthoredTarget::DeliveryPlan(id) => *id.as_bytes(), 1223 }, 1224 } 1225 } 1226 1227 fn command_phase(command: &AuthoredAtomicCommand) -> &'static str { 1228 match command { 1229 AuthoredAtomicCommand::Prepare(_) | AuthoredAtomicCommand::PrepareFromDraft(_) => "prepare", 1230 AuthoredAtomicCommand::Claim(_) => "claim", 1231 AuthoredAtomicCommand::ApplySigned(_) | AuthoredAtomicCommand::RecordSigned(_) => "signing", 1232 AuthoredAtomicCommand::ApplyAdmission(_) => "admission", 1233 AuthoredAtomicCommand::ApplyDelivery(_) 1234 | AuthoredAtomicCommand::RecordDelivery(_) 1235 | AuthoredAtomicCommand::ReconcileDelivery(_) => "delivery", 1236 AuthoredAtomicCommand::ApplyFailure(value) => match value.failure().phase() { 1237 WorkPhase::Signing => "signing_failure", 1238 WorkPhase::Admission => "admission_failure", 1239 WorkPhase::Delivery => "delivery_failure", 1240 }, 1241 AuthoredAtomicCommand::Cancel(_) => "cancel", 1242 } 1243 } 1244 1245 const fn origin_name(value: ArtifactOrigin) -> &'static str { 1246 match value { 1247 ArtifactOrigin::Planned => "planned", 1248 ArtifactOrigin::ImportedSigned => "imported_signed", 1249 } 1250 } 1251 1252 const fn signing_name(value: SigningState) -> &'static str { 1253 match value { 1254 SigningState::Planned => "planned", 1255 SigningState::Signed => "signed", 1256 SigningState::Retryable => "retryable", 1257 SigningState::Indeterminate => "indeterminate", 1258 SigningState::FailedTerminal => "failed_terminal", 1259 SigningState::Cancelled => "cancelled", 1260 } 1261 } 1262 1263 const fn physical_signing_name(artifact: &AuthoredArtifact) -> &'static str { 1264 if artifact.signed().is_some() { 1265 "signed" 1266 } else { 1267 signing_name(artifact.signing_state()) 1268 } 1269 } 1270 1271 const fn signing_stop(artifact: &AuthoredArtifact) -> Option<&'static str> { 1272 if artifact.signed().is_some() 1273 && matches!( 1274 artifact.signing_state(), 1275 SigningState::Cancelled | SigningState::FailedTerminal 1276 ) 1277 { 1278 Some(signing_name(artifact.signing_state())) 1279 } else { 1280 None 1281 } 1282 } 1283 1284 const fn admission_name(value: AdmissionState) -> &'static str { 1285 match value { 1286 AdmissionState::Pending => "pending", 1287 AdmissionState::Inserted => "inserted", 1288 AdmissionState::Duplicate => "duplicate", 1289 AdmissionState::Retryable => "retryable", 1290 AdmissionState::Rejected => "rejected", 1291 AdmissionState::Cancelled => "cancelled", 1292 } 1293 } 1294 1295 const fn delivery_state_name(value: AuthoredDeliveryState) -> &'static str { 1296 match value { 1297 AuthoredDeliveryState::Pending => "pending", 1298 AuthoredDeliveryState::Retryable => "retryable", 1299 AuthoredDeliveryState::Satisfied => "satisfied", 1300 AuthoredDeliveryState::Exhausted => "exhausted", 1301 AuthoredDeliveryState::FailedTerminal => "failed_terminal", 1302 AuthoredDeliveryState::Cancelled => "cancelled", 1303 } 1304 } 1305 1306 const fn satisfaction_name(value: SatisfactionState) -> &'static str { 1307 match value { 1308 SatisfactionState::Satisfied => "satisfied", 1309 SatisfactionState::Pending => "pending", 1310 SatisfactionState::Exhausted => "exhausted", 1311 } 1312 } 1313 1314 fn array<const N: usize>(bytes: Vec<u8>) -> Result<[u8; N], Error> { 1315 bytes.try_into().map_err(|_| Error::AtomicCommitFailed) 1316 } 1317 1318 fn i64_from_u64(value: u64) -> Result<i64, Error> { 1319 i64::try_from(value).map_err(|_| Error::AtomicCommitFailed) 1320 } 1321 1322 fn u64_from_i64(value: i64) -> Result<u64, Error> { 1323 u64::try_from(value).map_err(|_| Error::AtomicCommitFailed) 1324 } 1325 1326 fn column<T>(row: &SqliteRow, name: &str) -> Result<T, Error> 1327 where 1328 for<'decode> T: sqlx::Decode<'decode, Sqlite> + sqlx::Type<Sqlite>, 1329 { 1330 row.try_get(name).map_err(|_| Error::AtomicCommitFailed) 1331 } 1332 1333 #[cfg(test)] 1334 #[cfg_attr(coverage_nightly, coverage(off))] 1335 mod tests { 1336 use super::*; 1337 use crate::migration::runtime::{MIGRATIONS, migration_sql}; 1338 use core::num::{NonZeroU32, NonZeroU64}; 1339 use radroots_event::{GenericEventDraft, SignedEvent, wire::v1::Nip01EventWire}; 1340 use radroots_event_codec::authoring::AuthoredEventPlan; 1341 use radroots_storage::{ 1342 authored::{ 1343 AdmissionState, AuthoredArtifact, FailureClass, RetrySchedule, WorkClaim, WorkFailure, 1344 WorkPhase, 1345 }, 1346 authored_atomic::{ 1347 ApplyAdmissionResult, ApplyDeliveryAttempt, ApplySignedArtifact, ApplyWorkFailure, 1348 AuthoredWorkTarget, CancelAuthoredTarget, CancelAuthoredWork, ClaimAuthoredWork, 1349 PrepareAuthoredOperation, 1350 }, 1351 authored_delivery::{AuthoredDeliveryIntent, AuthoredDeliveryPlan}, 1352 event::SourceGeneration, 1353 status::EventStoreMode, 1354 }; 1355 use radroots_transport::{ 1356 DeliveryReceipt, SinkFailure, Target, TargetSet, 1357 outcome::{DeliveryOutcome, Retryability}, 1358 policy::{SatisfactionClass, SatisfactionPolicy, TargetPolicy}, 1359 sink::DeliveryTargetReceipt, 1360 }; 1361 use sqlx::sqlite::SqlitePoolOptions; 1362 1363 const AUTHOR: &str = "585591529da0bab31b3b1b1f986611cf5f435dca84f978c89ee8a40cca7103df"; 1364 1365 async fn store(mode: EventStoreMode) -> SqliteStorage { 1366 let generation = SourceGeneration::new([91; 32]).expect("generation"); 1367 let pool = SqlitePoolOptions::new() 1368 .max_connections(1) 1369 .connect("sqlite::memory:") 1370 .await 1371 .expect("memory SQLite"); 1372 sqlx::query("PRAGMA foreign_keys = ON") 1373 .execute(&pool) 1374 .await 1375 .expect("foreign keys"); 1376 for migration in MIGRATIONS { 1377 sqlx::raw_sql(migration_sql(migration.version()).expect("registered SQL")) 1378 .execute(&pool) 1379 .await 1380 .expect("runtime migration"); 1381 } 1382 sqlx::query( 1383 "INSERT INTO radroots_runtime_source_generations ( 1384 generation, sequence_head, state, created_at_unix_ms, retired_at_unix_ms 1385 ) VALUES (?, 0, 'active', 1, NULL)", 1386 ) 1387 .bind(generation.as_bytes().as_slice()) 1388 .execute(&pool) 1389 .await 1390 .expect("source generation"); 1391 SqliteStorage::new(pool, generation, mode) 1392 } 1393 1394 fn plan() -> AuthoredEventPlan { 1395 AuthoredEventPlan::from_generic( 1396 GenericEventDraft::new( 1397 "radroots.social.geochat.v1", 1398 20_000, 1399 1_800_100_001, 1400 Vec::new(), 1401 "sqlite authored operation", 1402 AUTHOR, 1403 ) 1404 .expect("generic draft"), 1405 ) 1406 .expect("authored plan") 1407 } 1408 1409 pub(super) fn signed(plan: &AuthoredEventPlan) -> SignedEvent { 1410 let wire = Nip01EventWire { 1411 id: plan.expected_event_id().to_hex(), 1412 pubkey: plan.author().to_hex(), 1413 created_at: plan.created_at(), 1414 kind: plan.body().kind(), 1415 tags: plan.body().tags().to_vec(), 1416 content: plan.body().content().to_owned(), 1417 sig: "44".repeat(64), 1418 extra: Default::default(), 1419 }; 1420 let raw = serde_json::to_string(&wire).expect("raw event"); 1421 SignedEvent::from_wire_verified_id(wire, raw).expect("signed event") 1422 } 1423 1424 pub(super) fn ids() -> ( 1425 OperationInstanceId, 1426 AuthoredArtifactId, 1427 AuthoredDeliveryPlanId, 1428 ) { 1429 ( 1430 OperationInstanceId::new([1; 16]).expect("operation"), 1431 AuthoredArtifactId::new([2; 16]).expect("artifact"), 1432 AuthoredDeliveryPlanId::new([3; 16]).expect("plan"), 1433 ) 1434 } 1435 1436 pub(super) fn prepare() -> (AuthoredAtomicCommand, AuthoredEventPlan) { 1437 let (operation_id, artifact_id, plan_id) = ids(); 1438 let event_plan = plan(); 1439 let artifact = AuthoredArtifact::planned(artifact_id, operation_id, 0, &event_plan, 10) 1440 .expect("artifact"); 1441 let operation = 1442 AuthoredOperation::new(operation_id, vec![artifact_id], 10).expect("operation"); 1443 let targets = TargetSet::new(vec![ 1444 Target::nostr_relay("wss://one.sqlite.example").expect("first"), 1445 Target::nostr_relay("wss://two.sqlite.example").expect("second"), 1446 ]) 1447 .expect("targets"); 1448 let intent = AuthoredDeliveryIntent::new( 1449 "sqlite-authored-delivery", 1450 targets, 1451 SatisfactionPolicy::new(SatisfactionClass::Accepted, TargetPolicy::any()), 1452 10_000, 1453 ) 1454 .expect("intent"); 1455 let delivery = 1456 AuthoredDeliveryPlan::new(plan_id, artifact_id, intent, 10).expect("delivery plan"); 1457 let command = PrepareAuthoredOperation::new( 1458 operation, 1459 vec![artifact], 1460 vec![delivery], 1461 AtomicCommitDigest::new([7; 32]), 1462 10, 1463 ) 1464 .expect("prepare"); 1465 (AuthoredAtomicCommand::Prepare(command), event_plan) 1466 } 1467 1468 pub(super) fn fence(claim: &WorkClaim) -> WorkFence { 1469 WorkFence::new(*claim.token(), claim.generation(), claim.row_revision()).unwrap() 1470 } 1471 1472 async fn prepared_store() -> (SqliteStorage, AuthoredEventPlan) { 1473 let store = store(EventStoreMode::ReadWrite).await; 1474 let (command, plan) = prepare(); 1475 store.execute_authored(command).await.unwrap(); 1476 (store, plan) 1477 } 1478 1479 async fn signed_store() -> SqliteStorage { 1480 let (store, event_plan) = prepared_store().await; 1481 let artifact = store.authored_artifact(ids().1).await.unwrap().unwrap(); 1482 let claim = WorkClaim::new( 1483 [4; 16], 1484 "sqlite-signer", 1485 NonZeroU64::MIN, 1486 11, 1487 50, 1488 artifact.revision(), 1489 ) 1490 .unwrap(); 1491 store 1492 .execute_authored(AuthoredAtomicCommand::Claim(ClaimAuthoredWork::new( 1493 ClaimAuthoredTarget::ArtifactSigning(ids().1), 1494 claim.clone(), 1495 ))) 1496 .await 1497 .unwrap(); 1498 store 1499 .execute_authored(AuthoredAtomicCommand::ApplySigned( 1500 ApplySignedArtifact::new(ids().1, fence(&claim), signed(&event_plan), 12).unwrap(), 1501 )) 1502 .await 1503 .unwrap(); 1504 store 1505 } 1506 1507 async fn claim_admission(store: &SqliteStorage) -> WorkClaim { 1508 let artifact = store.authored_artifact(ids().1).await.unwrap().unwrap(); 1509 let claim = WorkClaim::new( 1510 [5; 16], 1511 "sqlite-admission", 1512 NonZeroU64::new(2).unwrap(), 1513 13, 1514 50, 1515 artifact.revision(), 1516 ) 1517 .unwrap(); 1518 store 1519 .execute_authored(AuthoredAtomicCommand::Claim(ClaimAuthoredWork::new( 1520 ClaimAuthoredTarget::ArtifactAdmission(ids().1), 1521 claim.clone(), 1522 ))) 1523 .await 1524 .unwrap(); 1525 claim 1526 } 1527 1528 async fn claim_delivery(store: &SqliteStorage) -> WorkClaim { 1529 let plan = store 1530 .authored_delivery_plan(ids().2) 1531 .await 1532 .unwrap() 1533 .unwrap(); 1534 let claim = WorkClaim::new( 1535 [6; 16], 1536 "sqlite-delivery", 1537 NonZeroU64::new(3).unwrap(), 1538 13, 1539 50, 1540 plan.revision(), 1541 ) 1542 .unwrap(); 1543 store 1544 .execute_authored(AuthoredAtomicCommand::Claim(ClaimAuthoredWork::new( 1545 ClaimAuthoredTarget::DeliveryPlan(ids().2), 1546 claim.clone(), 1547 ))) 1548 .await 1549 .unwrap(); 1550 claim 1551 } 1552 1553 async fn satisfied_store() -> SqliteStorage { 1554 let store = signed_store().await; 1555 let delivery_claim = claim_delivery(&store).await; 1556 let plan = store 1557 .authored_delivery_plan(ids().2) 1558 .await 1559 .unwrap() 1560 .unwrap(); 1561 let request = plan.request().unwrap(); 1562 let receipt = DeliveryReceipt::for_request( 1563 request, 1564 request 1565 .target_set() 1566 .targets() 1567 .iter() 1568 .cloned() 1569 .map(|target| DeliveryTargetReceipt::attempted(target, DeliveryOutcome::accepted())) 1570 .collect(), 1571 ) 1572 .unwrap(); 1573 store 1574 .execute_authored(AuthoredAtomicCommand::ApplyDelivery( 1575 ApplyDeliveryAttempt::new( 1576 ids().2, 1577 fence(&delivery_claim), 1578 DeliveryAttemptOutcome::Receipt(receipt), 1579 None, 1580 14, 1581 ) 1582 .unwrap(), 1583 )) 1584 .await 1585 .unwrap(); 1586 store 1587 } 1588 1589 #[tokio::test] 1590 async fn authored_workflows_match_memory_replay_and_exact_binding_semantics() { 1591 let store = store(EventStoreMode::ReadWrite).await; 1592 let (prepare, event_plan) = prepare(); 1593 let committed = store 1594 .execute_authored(prepare.clone()) 1595 .await 1596 .expect("prepare"); 1597 assert_eq!(committed.disposition(), AtomicCommitDisposition::Committed); 1598 assert_eq!( 1599 store 1600 .execute_authored(prepare.clone()) 1601 .await 1602 .expect("replay") 1603 .disposition(), 1604 AtomicCommitDisposition::Replay 1605 ); 1606 assert_eq!( 1607 store 1608 .authored_receipt(prepare.commit_id()) 1609 .await 1610 .expect("receipt") 1611 .expect("stored receipt") 1612 .outcome(), 1613 committed.outcome() 1614 ); 1615 1616 let artifact = store 1617 .authored_artifact(ids().1) 1618 .await 1619 .expect("artifact query") 1620 .expect("artifact"); 1621 let claim = WorkClaim::new( 1622 [4; 16], 1623 "sqlite-signer", 1624 NonZeroU64::MIN, 1625 11, 1626 50, 1627 artifact.revision(), 1628 ) 1629 .expect("claim"); 1630 store 1631 .execute_authored(AuthoredAtomicCommand::Claim(ClaimAuthoredWork::new( 1632 ClaimAuthoredTarget::ArtifactSigning(ids().1), 1633 claim.clone(), 1634 ))) 1635 .await 1636 .expect("claim signing"); 1637 store 1638 .execute_authored(AuthoredAtomicCommand::ApplySigned( 1639 ApplySignedArtifact::new( 1640 ids().1, 1641 WorkFence::new(*claim.token(), claim.generation(), claim.row_revision()) 1642 .expect("fence"), 1643 signed(&event_plan), 1644 12, 1645 ) 1646 .expect("apply signed"), 1647 )) 1648 .await 1649 .expect("signed"); 1650 let artifact = store 1651 .authored_artifact(ids().1) 1652 .await 1653 .expect("artifact query") 1654 .expect("artifact"); 1655 let delivery = store 1656 .authored_delivery_plan(ids().2) 1657 .await 1658 .expect("delivery query") 1659 .expect("delivery"); 1660 assert_eq!(artifact.signing_state(), SigningState::Signed); 1661 assert_eq!( 1662 delivery 1663 .request() 1664 .expect("bound request") 1665 .payload() 1666 .event() 1667 .raw_json(), 1668 artifact 1669 .signed() 1670 .expect("signed artifact") 1671 .event() 1672 .raw_json() 1673 ); 1674 assert_eq!( 1675 sqlx::query_scalar::<_, i64>( 1676 "SELECT COUNT(*) FROM radroots_runtime_authored_delivery_targets", 1677 ) 1678 .fetch_one(&store.pool) 1679 .await 1680 .expect("target count"), 1681 2 1682 ); 1683 } 1684 1685 #[tokio::test] 1686 async fn statement_failure_rolls_back_complete_preparation_and_receipt() { 1687 let store = store(EventStoreMode::ReadWrite).await; 1688 sqlx::query( 1689 "CREATE TEMP TRIGGER authored_target_fault 1690 BEFORE INSERT ON radroots_runtime_authored_delivery_targets 1691 BEGIN SELECT RAISE(ABORT, 'authored target fault'); END", 1692 ) 1693 .execute(&store.pool) 1694 .await 1695 .expect("fault trigger"); 1696 let (prepare, _) = prepare(); 1697 assert_eq!( 1698 store.execute_authored(prepare).await, 1699 Err(Error::BackendUnavailable) 1700 ); 1701 for (table, query) in [ 1702 ( 1703 "radroots_runtime_authored_operations", 1704 "SELECT COUNT(*) FROM radroots_runtime_authored_operations", 1705 ), 1706 ( 1707 "radroots_runtime_authored_artifacts", 1708 "SELECT COUNT(*) FROM radroots_runtime_authored_artifacts", 1709 ), 1710 ( 1711 "radroots_runtime_authored_delivery_plans", 1712 "SELECT COUNT(*) FROM radroots_runtime_authored_delivery_plans", 1713 ), 1714 ( 1715 "radroots_runtime_authored_delivery_targets", 1716 "SELECT COUNT(*) FROM radroots_runtime_authored_delivery_targets", 1717 ), 1718 ( 1719 "radroots_runtime_authored_atomic_commits", 1720 "SELECT COUNT(*) FROM radroots_runtime_authored_atomic_commits", 1721 ), 1722 ] { 1723 assert_eq!( 1724 sqlx::query_scalar::<_, i64>(query) 1725 .fetch_one(&store.pool) 1726 .await 1727 .expect("row count"), 1728 0, 1729 "{table} retained partial state" 1730 ); 1731 } 1732 } 1733 1734 #[tokio::test] 1735 async fn corrupt_snapshots_children_claims_and_oversized_wires_fail_closed() { 1736 let store = store(EventStoreMode::ReadWrite).await; 1737 let (prepare, _) = prepare(); 1738 store.execute_authored(prepare).await.expect("prepare"); 1739 1740 assert!( 1741 sqlx::query( 1742 "UPDATE radroots_runtime_authored_delivery_plans 1743 SET claim_token = zeroblob(16) WHERE plan_id = ?", 1744 ) 1745 .bind(ids().2.as_bytes().as_slice()) 1746 .execute(&store.pool) 1747 .await 1748 .is_err() 1749 ); 1750 assert!( 1751 sqlx::query( 1752 "UPDATE radroots_runtime_authored_artifacts 1753 SET snapshot = zeroblob(4194305) WHERE artifact_id = ?", 1754 ) 1755 .bind(ids().1.as_bytes().as_slice()) 1756 .execute(&store.pool) 1757 .await 1758 .is_err() 1759 ); 1760 sqlx::query( 1761 "UPDATE radroots_runtime_authored_delivery_targets 1762 SET target_fingerprint = 'forged' WHERE plan_id = ? AND ordinal = 0", 1763 ) 1764 .bind(ids().2.as_bytes().as_slice()) 1765 .execute(&store.pool) 1766 .await 1767 .expect("forge child row"); 1768 assert_eq!( 1769 store.authored_delivery_plan(ids().2).await, 1770 Err(Error::InvalidAuthoredDeliveryPlan) 1771 ); 1772 sqlx::query( 1773 "UPDATE radroots_runtime_authored_artifacts SET snapshot = x'7b7d' 1774 WHERE artifact_id = ?", 1775 ) 1776 .bind(ids().1.as_bytes().as_slice()) 1777 .execute(&store.pool) 1778 .await 1779 .expect("forge snapshot"); 1780 assert_eq!( 1781 store.authored_artifact(ids().1).await, 1782 Err(Error::AtomicCommitFailed) 1783 ); 1784 } 1785 1786 #[tokio::test] 1787 async fn authored_admission_and_delivery_commands_round_trip_exact_state() { 1788 let store = signed_store().await; 1789 let admission_claim = claim_admission(&store).await; 1790 store 1791 .execute_authored(AuthoredAtomicCommand::ApplyAdmission( 1792 ApplyAdmissionResult::new( 1793 ids().1, 1794 fence(&admission_claim), 1795 AdmissionState::Inserted, 1796 None, 1797 None, 1798 14, 1799 ) 1800 .unwrap(), 1801 )) 1802 .await 1803 .unwrap(); 1804 assert_eq!( 1805 store 1806 .authored_artifact(ids().1) 1807 .await 1808 .unwrap() 1809 .unwrap() 1810 .admission_state(), 1811 AdmissionState::Inserted 1812 ); 1813 1814 let delivery_claim = claim_delivery(&store).await; 1815 let plan = store 1816 .authored_delivery_plan(ids().2) 1817 .await 1818 .unwrap() 1819 .unwrap(); 1820 let request = plan.request().unwrap(); 1821 let receipt = DeliveryReceipt::for_request( 1822 request, 1823 request 1824 .target_set() 1825 .targets() 1826 .iter() 1827 .cloned() 1828 .map(|target| DeliveryTargetReceipt::attempted(target, DeliveryOutcome::accepted())) 1829 .collect(), 1830 ) 1831 .unwrap(); 1832 store 1833 .execute_authored(AuthoredAtomicCommand::ApplyDelivery( 1834 ApplyDeliveryAttempt::new( 1835 ids().2, 1836 fence(&delivery_claim), 1837 DeliveryAttemptOutcome::Receipt(receipt), 1838 None, 1839 14, 1840 ) 1841 .unwrap(), 1842 )) 1843 .await 1844 .unwrap(); 1845 assert_eq!( 1846 store 1847 .authored_delivery_plan(ids().2) 1848 .await 1849 .unwrap() 1850 .unwrap() 1851 .state(), 1852 AuthoredDeliveryState::Satisfied 1853 ); 1854 } 1855 1856 #[tokio::test] 1857 async fn authored_failure_commands_persist_each_phase_and_retry_class() { 1858 let (store, _) = prepared_store().await; 1859 let artifact = store.authored_artifact(ids().1).await.unwrap().unwrap(); 1860 let claim = WorkClaim::new( 1861 [4; 16], 1862 "sqlite-signer", 1863 NonZeroU64::MIN, 1864 11, 1865 50, 1866 artifact.revision(), 1867 ) 1868 .unwrap(); 1869 store 1870 .execute_authored(AuthoredAtomicCommand::Claim(ClaimAuthoredWork::new( 1871 ClaimAuthoredTarget::ArtifactSigning(ids().1), 1872 claim.clone(), 1873 ))) 1874 .await 1875 .unwrap(); 1876 let failure = WorkFailure::new( 1877 "terminal_signer", 1878 WorkPhase::Signing, 1879 FailureClass::Terminal, 1880 None, 1881 None, 1882 ) 1883 .unwrap(); 1884 store 1885 .execute_authored(AuthoredAtomicCommand::ApplyFailure( 1886 ApplyWorkFailure::new( 1887 AuthoredWorkTarget::Artifact(ids().1), 1888 fence(&claim), 1889 failure, 1890 None, 1891 12, 1892 ) 1893 .unwrap(), 1894 )) 1895 .await 1896 .unwrap(); 1897 assert_eq!( 1898 store 1899 .authored_artifact(ids().1) 1900 .await 1901 .unwrap() 1902 .unwrap() 1903 .signing_state(), 1904 SigningState::FailedTerminal 1905 ); 1906 1907 let store = signed_store().await; 1908 let claim = claim_admission(&store).await; 1909 let failure = WorkFailure::new( 1910 "temporary_admission", 1911 WorkPhase::Admission, 1912 FailureClass::Retryable, 1913 Some(30), 1914 None, 1915 ) 1916 .unwrap(); 1917 let retry = RetrySchedule::new(NonZeroU32::MIN, 30, failure.clone()).unwrap(); 1918 store 1919 .execute_authored(AuthoredAtomicCommand::ApplyFailure( 1920 ApplyWorkFailure::new( 1921 AuthoredWorkTarget::Artifact(ids().1), 1922 fence(&claim), 1923 failure, 1924 Some(retry), 1925 14, 1926 ) 1927 .unwrap(), 1928 )) 1929 .await 1930 .unwrap(); 1931 assert_eq!( 1932 store 1933 .authored_artifact(ids().1) 1934 .await 1935 .unwrap() 1936 .unwrap() 1937 .admission_state(), 1938 AdmissionState::Retryable 1939 ); 1940 1941 let store = signed_store().await; 1942 let claim = claim_delivery(&store).await; 1943 let failure = WorkFailure::new( 1944 "relay_terminal", 1945 WorkPhase::Delivery, 1946 FailureClass::Terminal, 1947 None, 1948 Some("terminal".to_owned()), 1949 ) 1950 .unwrap(); 1951 store 1952 .execute_authored(AuthoredAtomicCommand::ApplyFailure( 1953 ApplyWorkFailure::new( 1954 AuthoredWorkTarget::DeliveryPlan(ids().2), 1955 fence(&claim), 1956 failure, 1957 None, 1958 14, 1959 ) 1960 .unwrap(), 1961 )) 1962 .await 1963 .unwrap(); 1964 assert_eq!( 1965 store 1966 .authored_delivery_plan(ids().2) 1967 .await 1968 .unwrap() 1969 .unwrap() 1970 .state(), 1971 AuthoredDeliveryState::FailedTerminal 1972 ); 1973 } 1974 1975 #[tokio::test] 1976 async fn authored_sink_failure_and_every_cancellation_target_are_durable() { 1977 let store = signed_store().await; 1978 let claim = claim_delivery(&store).await; 1979 let plan = store 1980 .authored_delivery_plan(ids().2) 1981 .await 1982 .unwrap() 1983 .unwrap(); 1984 let failure = SinkFailure::for_request( 1985 plan.request().unwrap(), 1986 "relay_unavailable", 1987 Retryability::Retryable, 1988 Some(30), 1989 None, 1990 Vec::new(), 1991 ) 1992 .unwrap(); 1993 let work_failure = WorkFailure::new( 1994 "relay_unavailable", 1995 WorkPhase::Delivery, 1996 FailureClass::Retryable, 1997 Some(30), 1998 None, 1999 ) 2000 .unwrap(); 2001 let retry = RetrySchedule::new(NonZeroU32::MIN, 30, work_failure).unwrap(); 2002 store 2003 .execute_authored(AuthoredAtomicCommand::ApplyDelivery( 2004 ApplyDeliveryAttempt::new( 2005 ids().2, 2006 fence(&claim), 2007 DeliveryAttemptOutcome::SinkFailure(failure), 2008 Some(retry), 2009 14, 2010 ) 2011 .unwrap(), 2012 )) 2013 .await 2014 .unwrap(); 2015 assert_eq!( 2016 store 2017 .authored_delivery_plan(ids().2) 2018 .await 2019 .unwrap() 2020 .unwrap() 2021 .state(), 2022 AuthoredDeliveryState::Retryable 2023 ); 2024 2025 let (store, _) = prepared_store().await; 2026 let artifact = store.authored_artifact(ids().1).await.unwrap().unwrap(); 2027 store 2028 .execute_authored(AuthoredAtomicCommand::Cancel( 2029 CancelAuthoredWork::new( 2030 CancelAuthoredTarget::ArtifactSigning(ids().1), 2031 artifact.revision(), 2032 11, 2033 ) 2034 .unwrap(), 2035 )) 2036 .await 2037 .unwrap(); 2038 2039 let store = signed_store().await; 2040 let artifact = store.authored_artifact(ids().1).await.unwrap().unwrap(); 2041 store 2042 .execute_authored(AuthoredAtomicCommand::Cancel( 2043 CancelAuthoredWork::new( 2044 CancelAuthoredTarget::ArtifactAdmission(ids().1), 2045 artifact.revision(), 2046 13, 2047 ) 2048 .unwrap(), 2049 )) 2050 .await 2051 .unwrap(); 2052 2053 let store = signed_store().await; 2054 let plan = store 2055 .authored_delivery_plan(ids().2) 2056 .await 2057 .unwrap() 2058 .unwrap(); 2059 store 2060 .execute_authored(AuthoredAtomicCommand::Cancel( 2061 CancelAuthoredWork::new( 2062 CancelAuthoredTarget::DeliveryPlan(ids().2), 2063 plan.revision(), 2064 13, 2065 ) 2066 .unwrap(), 2067 )) 2068 .await 2069 .unwrap(); 2070 } 2071 2072 #[tokio::test] 2073 async fn authored_missing_rows_conflicts_and_cross_table_integrity_fail_closed() { 2074 let store = store(EventStoreMode::ReadWrite).await; 2075 assert!(store.authored_operation(ids().0).await.unwrap().is_none()); 2076 assert!(store.authored_artifact(ids().1).await.unwrap().is_none()); 2077 assert!( 2078 store 2079 .authored_delivery_plan(ids().2) 2080 .await 2081 .unwrap() 2082 .is_none() 2083 ); 2084 assert!( 2085 store 2086 .authored_receipt(AtomicCommitId::new([9; 16]).unwrap()) 2087 .await 2088 .unwrap() 2089 .is_none() 2090 ); 2091 2092 let (command, _) = prepare(); 2093 store.execute_authored(command.clone()).await.unwrap(); 2094 assert_eq!( 2095 store.execute_authored(command).await.unwrap().disposition(), 2096 AtomicCommitDisposition::Replay 2097 ); 2098 let (conflict, _) = prepare(); 2099 let AuthoredAtomicCommand::Prepare(value) = conflict else { 2100 unreachable!() 2101 }; 2102 let conflict = PrepareAuthoredOperation::new( 2103 value.operation().clone(), 2104 value.artifacts().to_vec(), 2105 value.delivery_plans().to_vec(), 2106 AtomicCommitDigest::new([8; 32]), 2107 value.requested_at_unix_ms(), 2108 ) 2109 .unwrap(); 2110 assert_eq!( 2111 store 2112 .execute_authored(AuthoredAtomicCommand::Prepare(conflict)) 2113 .await, 2114 Err(Error::AtomicCommitConflict) 2115 ); 2116 2117 sqlx::query("DELETE FROM radroots_runtime_authored_artifacts WHERE artifact_id = ?") 2118 .bind(ids().1.as_bytes().as_slice()) 2119 .execute(&store.pool) 2120 .await 2121 .unwrap(); 2122 assert_eq!( 2123 store.authored_operation(ids().0).await, 2124 Err(Error::InvalidAuthoredOperation) 2125 ); 2126 } 2127 2128 #[tokio::test] 2129 async fn denormalized_authored_columns_are_verified_against_canonical_snapshots() { 2130 let (store, _) = prepared_store().await; 2131 2132 macro_rules! reject_operation_shadow { 2133 ($update:literal, $restore:literal) => {{ 2134 sqlx::query($update).execute(&store.pool).await.unwrap(); 2135 assert_eq!( 2136 store.authored_operation(ids().0).await, 2137 Err(Error::InvalidAuthoredOperation) 2138 ); 2139 sqlx::query($restore).execute(&store.pool).await.unwrap(); 2140 }}; 2141 } 2142 reject_operation_shadow!( 2143 "UPDATE radroots_runtime_authored_operations SET artifact_count = 2", 2144 "UPDATE radroots_runtime_authored_operations SET artifact_count = 1" 2145 ); 2146 reject_operation_shadow!( 2147 "UPDATE radroots_runtime_authored_operations SET created_at_unix_ms = 9", 2148 "UPDATE radroots_runtime_authored_operations SET created_at_unix_ms = 10" 2149 ); 2150 reject_operation_shadow!( 2151 "UPDATE radroots_runtime_authored_operations SET updated_at_unix_ms = 11", 2152 "UPDATE radroots_runtime_authored_operations SET updated_at_unix_ms = 10" 2153 ); 2154 reject_operation_shadow!( 2155 "UPDATE radroots_runtime_authored_operations SET revision = 2", 2156 "UPDATE radroots_runtime_authored_operations SET revision = 1" 2157 ); 2158 2159 macro_rules! reject_artifact_shadow { 2160 ($update:literal, $restore:literal) => {{ 2161 sqlx::query($update).execute(&store.pool).await.unwrap(); 2162 assert_eq!( 2163 store.authored_artifact(ids().1).await, 2164 Err(Error::InvalidAuthoredArtifact) 2165 ); 2166 sqlx::query($restore).execute(&store.pool).await.unwrap(); 2167 }}; 2168 } 2169 reject_artifact_shadow!( 2170 "UPDATE radroots_runtime_authored_artifacts SET ordinal = 1", 2171 "UPDATE radroots_runtime_authored_artifacts SET ordinal = 0" 2172 ); 2173 reject_artifact_shadow!( 2174 "UPDATE radroots_runtime_authored_artifacts SET signing_state = 'retryable'", 2175 "UPDATE radroots_runtime_authored_artifacts SET signing_state = 'planned'" 2176 ); 2177 reject_artifact_shadow!( 2178 "UPDATE radroots_runtime_authored_artifacts SET admission_state = 'inserted'", 2179 "UPDATE radroots_runtime_authored_artifacts SET admission_state = 'pending'" 2180 ); 2181 reject_artifact_shadow!( 2182 "UPDATE radroots_runtime_authored_artifacts SET retry_not_before_unix_ms = 30", 2183 "UPDATE radroots_runtime_authored_artifacts SET retry_not_before_unix_ms = NULL" 2184 ); 2185 reject_artifact_shadow!( 2186 "UPDATE radroots_runtime_authored_artifacts SET last_failure_code = 'forged'", 2187 "UPDATE radroots_runtime_authored_artifacts SET last_failure_code = NULL" 2188 ); 2189 reject_artifact_shadow!( 2190 "UPDATE radroots_runtime_authored_artifacts SET created_at_unix_ms = 9", 2191 "UPDATE radroots_runtime_authored_artifacts SET created_at_unix_ms = 10" 2192 ); 2193 reject_artifact_shadow!( 2194 "UPDATE radroots_runtime_authored_artifacts SET updated_at_unix_ms = 11", 2195 "UPDATE radroots_runtime_authored_artifacts SET updated_at_unix_ms = 10" 2196 ); 2197 reject_artifact_shadow!( 2198 "UPDATE radroots_runtime_authored_artifacts SET revision = 2", 2199 "UPDATE radroots_runtime_authored_artifacts SET revision = 1" 2200 ); 2201 2202 macro_rules! reject_plan_shadow { 2203 ($update:literal, $restore:literal) => {{ 2204 sqlx::query($update).execute(&store.pool).await.unwrap(); 2205 assert_eq!( 2206 store.authored_delivery_plan(ids().2).await, 2207 Err(Error::InvalidAuthoredDeliveryPlan) 2208 ); 2209 sqlx::query($restore).execute(&store.pool).await.unwrap(); 2210 }}; 2211 } 2212 reject_plan_shadow!( 2213 "UPDATE radroots_runtime_authored_delivery_plans SET attempt_count = 1", 2214 "UPDATE radroots_runtime_authored_delivery_plans SET attempt_count = 0" 2215 ); 2216 reject_plan_shadow!( 2217 "UPDATE radroots_runtime_authored_delivery_plans SET last_failure_code = 'forged'", 2218 "UPDATE radroots_runtime_authored_delivery_plans SET last_failure_code = NULL" 2219 ); 2220 reject_plan_shadow!( 2221 "UPDATE radroots_runtime_authored_delivery_plans SET created_at_unix_ms = 9", 2222 "UPDATE radroots_runtime_authored_delivery_plans SET created_at_unix_ms = 10" 2223 ); 2224 reject_plan_shadow!( 2225 "UPDATE radroots_runtime_authored_delivery_plans SET updated_at_unix_ms = 11", 2226 "UPDATE radroots_runtime_authored_delivery_plans SET updated_at_unix_ms = 10" 2227 ); 2228 reject_plan_shadow!( 2229 "UPDATE radroots_runtime_authored_delivery_plans SET revision = 2", 2230 "UPDATE radroots_runtime_authored_delivery_plans SET revision = 1" 2231 ); 2232 reject_plan_shadow!( 2233 "UPDATE radroots_runtime_authored_delivery_targets SET ordinal = 9 WHERE ordinal = 0", 2234 "UPDATE radroots_runtime_authored_delivery_targets SET ordinal = 0 WHERE ordinal = 9" 2235 ); 2236 } 2237 2238 #[tokio::test] 2239 async fn authored_binary_claim_attempt_and_receipt_shadows_fail_closed() { 2240 assert_eq!( 2241 decode_snapshot::<serde_json::Value>(Vec::new()), 2242 Err(Error::AtomicCommitFailed) 2243 ); 2244 assert_eq!( 2245 decode_snapshot::<serde_json::Value>(vec![0; SNAPSHOT_MAX_BYTES + 1]), 2246 Err(Error::AtomicCommitFailed) 2247 ); 2248 assert_eq!( 2249 encode_snapshot(&"x".repeat(SNAPSHOT_MAX_BYTES + 1)), 2250 Err(Error::AtomicCommitFailed) 2251 ); 2252 2253 let (store, _) = prepared_store().await; 2254 let artifact = store.authored_artifact(ids().1).await.unwrap().unwrap(); 2255 let plan_wire = artifact.plan().unwrap().wire_json().to_vec(); 2256 sqlx::query("UPDATE radroots_runtime_authored_artifacts SET plan_wire = x'7b7d'") 2257 .execute(&store.pool) 2258 .await 2259 .unwrap(); 2260 assert_eq!( 2261 store.authored_artifact(ids().1).await, 2262 Err(Error::InvalidAuthoredArtifact) 2263 ); 2264 sqlx::query("UPDATE radroots_runtime_authored_artifacts SET plan_wire = ?") 2265 .bind(plan_wire.as_slice()) 2266 .execute(&store.pool) 2267 .await 2268 .unwrap(); 2269 sqlx::query( 2270 "UPDATE radroots_runtime_authored_artifacts 2271 SET origin = 'imported_signed', plan_wire = NULL, 2272 signing_state = 'signed', signed_raw_json = x'7b7d', signed_raw_sha256 = zeroblob(32)", 2273 ) 2274 .execute(&store.pool) 2275 .await 2276 .unwrap(); 2277 assert_eq!( 2278 store.authored_artifact(ids().1).await, 2279 Err(Error::InvalidAuthoredArtifact) 2280 ); 2281 2282 let (store, _) = prepared_store().await; 2283 let plan = store 2284 .authored_delivery_plan(ids().2) 2285 .await 2286 .unwrap() 2287 .unwrap(); 2288 sqlx::query( 2289 "UPDATE radroots_runtime_authored_delivery_plans SET request_digest = zeroblob(32)", 2290 ) 2291 .execute(&store.pool) 2292 .await 2293 .unwrap(); 2294 assert_eq!( 2295 store.authored_delivery_plan(ids().2).await, 2296 Err(Error::InvalidAuthoredDeliveryPlan) 2297 ); 2298 sqlx::query( 2299 "UPDATE radroots_runtime_authored_delivery_plans SET request_digest = ?, state = 'cancelled'", 2300 ) 2301 .bind(plan.request_digest().as_slice()) 2302 .execute(&store.pool) 2303 .await 2304 .unwrap(); 2305 assert_eq!( 2306 store.authored_delivery_plan(ids().2).await, 2307 Err(Error::InvalidAuthoredDeliveryPlan) 2308 ); 2309 sqlx::query("UPDATE radroots_runtime_authored_delivery_plans SET state = 'pending'") 2310 .execute(&store.pool) 2311 .await 2312 .unwrap(); 2313 let target = plan.intent().target_set().targets()[0].clone(); 2314 sqlx::query( 2315 "UPDATE radroots_runtime_authored_delivery_targets 2316 SET target_fingerprint = 'forged' WHERE ordinal = 0", 2317 ) 2318 .execute(&store.pool) 2319 .await 2320 .unwrap(); 2321 assert_eq!( 2322 store.authored_delivery_plan(ids().2).await, 2323 Err(Error::InvalidAuthoredDeliveryPlan) 2324 ); 2325 sqlx::query( 2326 "UPDATE radroots_runtime_authored_delivery_targets 2327 SET target_fingerprint = ?, target_snapshot = x'7b7d' WHERE ordinal = 0", 2328 ) 2329 .bind(target.fingerprint().as_str()) 2330 .execute(&store.pool) 2331 .await 2332 .unwrap(); 2333 assert_eq!( 2334 store.authored_delivery_plan(ids().2).await, 2335 Err(Error::AtomicCommitFailed) 2336 ); 2337 2338 let (store, _) = prepared_store().await; 2339 let artifact = store.authored_artifact(ids().1).await.unwrap().unwrap(); 2340 let claim = WorkClaim::new( 2341 [4; 16], 2342 "sqlite-signer", 2343 NonZeroU64::MIN, 2344 11, 2345 50, 2346 artifact.revision(), 2347 ) 2348 .unwrap(); 2349 store 2350 .execute_authored(AuthoredAtomicCommand::Claim(ClaimAuthoredWork::new( 2351 ClaimAuthoredTarget::ArtifactSigning(ids().1), 2352 claim, 2353 ))) 2354 .await 2355 .unwrap(); 2356 for (update, restore) in [ 2357 ( 2358 "UPDATE radroots_runtime_authored_artifacts SET signing_claim_token = zeroblob(16)", 2359 "UPDATE radroots_runtime_authored_artifacts SET signing_claim_token = x'04040404040404040404040404040404'", 2360 ), 2361 ( 2362 "UPDATE radroots_runtime_authored_artifacts SET signing_claim_generation = 2", 2363 "UPDATE radroots_runtime_authored_artifacts SET signing_claim_generation = 1", 2364 ), 2365 ( 2366 "UPDATE radroots_runtime_authored_artifacts SET signing_claim_revision = 2", 2367 "UPDATE radroots_runtime_authored_artifacts SET signing_claim_revision = 1", 2368 ), 2369 ( 2370 "UPDATE radroots_runtime_authored_artifacts SET signing_claim_expires_at_unix_ms = 51", 2371 "UPDATE radroots_runtime_authored_artifacts SET signing_claim_expires_at_unix_ms = 50", 2372 ), 2373 ] { 2374 sqlx::query(update).execute(&store.pool).await.unwrap(); 2375 assert_eq!( 2376 store.authored_artifact(ids().1).await, 2377 Err(Error::InvalidAuthoredArtifact) 2378 ); 2379 sqlx::query(restore).execute(&store.pool).await.unwrap(); 2380 } 2381 2382 let store = signed_store().await; 2383 claim_delivery(&store).await; 2384 for (update, restore) in [ 2385 ( 2386 "UPDATE radroots_runtime_authored_delivery_plans SET claim_token = zeroblob(16)", 2387 "UPDATE radroots_runtime_authored_delivery_plans SET claim_token = x'06060606060606060606060606060606'", 2388 ), 2389 ( 2390 "UPDATE radroots_runtime_authored_delivery_plans SET claim_generation = 4", 2391 "UPDATE radroots_runtime_authored_delivery_plans SET claim_generation = 3", 2392 ), 2393 ( 2394 "UPDATE radroots_runtime_authored_delivery_plans SET claim_revision = 99", 2395 "UPDATE radroots_runtime_authored_delivery_plans SET claim_revision = 2", 2396 ), 2397 ( 2398 "UPDATE radroots_runtime_authored_delivery_plans SET claim_expires_at_unix_ms = 51", 2399 "UPDATE radroots_runtime_authored_delivery_plans SET claim_expires_at_unix_ms = 50", 2400 ), 2401 ] { 2402 sqlx::query(update).execute(&store.pool).await.unwrap(); 2403 assert_eq!( 2404 store.authored_delivery_plan(ids().2).await, 2405 Err(Error::InvalidAuthoredDeliveryPlan) 2406 ); 2407 sqlx::query(restore).execute(&store.pool).await.unwrap(); 2408 } 2409 2410 let store = signed_store().await; 2411 let artifact = store.authored_artifact(ids().1).await.unwrap().unwrap(); 2412 let signed = artifact.signed().expect("signed artifact"); 2413 assert!( 2414 sqlx::query("UPDATE radroots_runtime_authored_artifacts SET signed_raw_json = x'7b7d'") 2415 .execute(&store.pool) 2416 .await 2417 .is_err() 2418 ); 2419 // Simulate corruption below the write guard to retain read-side validation coverage. 2420 sqlx::query("DROP TRIGGER radroots_runtime_authored_artifacts_signed_fact_guard") 2421 .execute(&store.pool) 2422 .await 2423 .unwrap(); 2424 sqlx::query("UPDATE radroots_runtime_authored_artifacts SET signed_raw_json = x'7b7d'") 2425 .execute(&store.pool) 2426 .await 2427 .unwrap(); 2428 assert_eq!( 2429 store.authored_artifact(ids().1).await, 2430 Err(Error::InvalidAuthoredArtifact) 2431 ); 2432 sqlx::query( 2433 "UPDATE radroots_runtime_authored_artifacts 2434 SET signed_raw_json = ?, signed_raw_sha256 = zeroblob(32)", 2435 ) 2436 .bind(signed.event().raw_json().as_bytes()) 2437 .execute(&store.pool) 2438 .await 2439 .unwrap(); 2440 assert_eq!( 2441 store.authored_artifact(ids().1).await, 2442 Err(Error::InvalidAuthoredArtifact) 2443 ); 2444 2445 let store = signed_store().await; 2446 claim_admission(&store).await; 2447 sqlx::query( 2448 "UPDATE radroots_runtime_authored_artifacts 2449 SET admission_claim_generation = admission_claim_generation + 1", 2450 ) 2451 .execute(&store.pool) 2452 .await 2453 .unwrap(); 2454 assert_eq!( 2455 store.authored_artifact(ids().1).await, 2456 Err(Error::InvalidAuthoredArtifact) 2457 ); 2458 2459 let (store, _) = prepared_store().await; 2460 sqlx::query("DELETE FROM radroots_runtime_authored_delivery_targets WHERE plan_id = ?") 2461 .bind(ids().2.as_bytes().as_slice()) 2462 .execute(&store.pool) 2463 .await 2464 .unwrap(); 2465 assert_eq!( 2466 store.authored_delivery_plan(ids().2).await, 2467 Err(Error::InvalidAuthoredDeliveryPlan) 2468 ); 2469 2470 let (store, _) = prepared_store().await; 2471 sqlx::query( 2472 "INSERT INTO radroots_runtime_authored_delivery_attempts ( 2473 plan_id, attempt, satisfaction, recorded_at_unix_ms, outcome_snapshot 2474 ) VALUES (?, 1, 'pending', 20, x'7b7d')", 2475 ) 2476 .bind(ids().2.as_bytes().as_slice()) 2477 .execute(&store.pool) 2478 .await 2479 .unwrap(); 2480 assert_eq!( 2481 store.authored_delivery_plan(ids().2).await, 2482 Err(Error::InvalidAuthoredDeliveryPlan) 2483 ); 2484 2485 let (store, _) = prepared_store().await; 2486 let other_target = Target::nostr_relay("wss://shadow.example").unwrap(); 2487 sqlx::query( 2488 "UPDATE radroots_runtime_authored_delivery_targets 2489 SET target_snapshot = ? WHERE plan_id = ? AND ordinal = 0", 2490 ) 2491 .bind(encode_snapshot(&other_target).unwrap()) 2492 .bind(ids().2.as_bytes().as_slice()) 2493 .execute(&store.pool) 2494 .await 2495 .unwrap(); 2496 assert_eq!( 2497 store.authored_delivery_plan(ids().2).await, 2498 Err(Error::InvalidAuthoredDeliveryPlan) 2499 ); 2500 2501 let store = satisfied_store().await; 2502 for (update, restore) in [ 2503 ( 2504 "UPDATE radroots_runtime_authored_delivery_attempts SET attempt = 2", 2505 "UPDATE radroots_runtime_authored_delivery_attempts SET attempt = 1", 2506 ), 2507 ( 2508 "UPDATE radroots_runtime_authored_delivery_attempts SET satisfaction = 'pending'", 2509 "UPDATE radroots_runtime_authored_delivery_attempts SET satisfaction = 'satisfied'", 2510 ), 2511 ( 2512 "UPDATE radroots_runtime_authored_delivery_attempts SET recorded_at_unix_ms = 15", 2513 "UPDATE radroots_runtime_authored_delivery_attempts SET recorded_at_unix_ms = 14", 2514 ), 2515 ] { 2516 sqlx::query(update).execute(&store.pool).await.unwrap(); 2517 assert_eq!( 2518 store.authored_delivery_plan(ids().2).await, 2519 Err(Error::InvalidAuthoredDeliveryPlan) 2520 ); 2521 sqlx::query(restore).execute(&store.pool).await.unwrap(); 2522 } 2523 let plan = store 2524 .authored_delivery_plan(ids().2) 2525 .await 2526 .unwrap() 2527 .unwrap(); 2528 let request = plan.request().unwrap(); 2529 let different = DeliveryAttemptOutcome::Receipt( 2530 DeliveryReceipt::for_request( 2531 request, 2532 request 2533 .target_set() 2534 .targets() 2535 .iter() 2536 .cloned() 2537 .map(|target| { 2538 DeliveryTargetReceipt::attempted(target, DeliveryOutcome::unavailable()) 2539 }) 2540 .collect(), 2541 ) 2542 .unwrap(), 2543 ); 2544 sqlx::query("UPDATE radroots_runtime_authored_delivery_attempts SET outcome_snapshot = ?") 2545 .bind(encode_snapshot(&different).unwrap()) 2546 .execute(&store.pool) 2547 .await 2548 .unwrap(); 2549 assert_eq!( 2550 store.authored_delivery_plan(ids().2).await, 2551 Err(Error::InvalidAuthoredDeliveryPlan) 2552 ); 2553 2554 async fn disable_foreign_keys(store: &SqliteStorage) { 2555 sqlx::query("PRAGMA foreign_keys = OFF") 2556 .execute(&store.pool) 2557 .await 2558 .unwrap(); 2559 } 2560 2561 let (store, _) = prepared_store().await; 2562 disable_foreign_keys(&store).await; 2563 sqlx::query("UPDATE radroots_runtime_authored_operations SET operation_id = ?") 2564 .bind([41; 16].as_slice()) 2565 .execute(&store.pool) 2566 .await 2567 .unwrap(); 2568 assert_eq!( 2569 store 2570 .authored_operation(OperationInstanceId::new([41; 16]).unwrap()) 2571 .await, 2572 Err(Error::InvalidAuthoredOperation) 2573 ); 2574 2575 let (store, _) = prepared_store().await; 2576 disable_foreign_keys(&store).await; 2577 sqlx::query("UPDATE radroots_runtime_authored_artifacts SET artifact_id = ?") 2578 .bind([42; 16].as_slice()) 2579 .execute(&store.pool) 2580 .await 2581 .unwrap(); 2582 assert_eq!( 2583 store 2584 .authored_artifact(AuthoredArtifactId::new([42; 16]).unwrap()) 2585 .await, 2586 Err(Error::InvalidAuthoredArtifact) 2587 ); 2588 2589 let (store, _) = prepared_store().await; 2590 disable_foreign_keys(&store).await; 2591 sqlx::query("UPDATE radroots_runtime_authored_artifacts SET operation_id = ?") 2592 .bind([43; 16].as_slice()) 2593 .execute(&store.pool) 2594 .await 2595 .unwrap(); 2596 assert_eq!( 2597 store.authored_artifact(ids().1).await, 2598 Err(Error::InvalidAuthoredArtifact) 2599 ); 2600 2601 let (store, _) = prepared_store().await; 2602 disable_foreign_keys(&store).await; 2603 sqlx::query("UPDATE radroots_runtime_authored_delivery_plans SET plan_id = ?") 2604 .bind([44; 16].as_slice()) 2605 .execute(&store.pool) 2606 .await 2607 .unwrap(); 2608 assert_eq!( 2609 store 2610 .authored_delivery_plan(AuthoredDeliveryPlanId::new([44; 16]).unwrap()) 2611 .await, 2612 Err(Error::InvalidAuthoredDeliveryPlan) 2613 ); 2614 2615 let (store, _) = prepared_store().await; 2616 disable_foreign_keys(&store).await; 2617 sqlx::query("UPDATE radroots_runtime_authored_delivery_plans SET artifact_id = ?") 2618 .bind([45; 16].as_slice()) 2619 .execute(&store.pool) 2620 .await 2621 .unwrap(); 2622 assert_eq!( 2623 store.authored_delivery_plan(ids().2).await, 2624 Err(Error::InvalidAuthoredDeliveryPlan) 2625 ); 2626 } 2627 2628 #[tokio::test] 2629 async fn read_only_and_ready_query_plan_contracts_are_enforced() { 2630 let store = store(EventStoreMode::ReadOnly).await; 2631 assert_eq!( 2632 store.execute_authored(prepare().0).await, 2633 Err(Error::BackendUnavailable) 2634 ); 2635 let plan = sqlx::query( 2636 "EXPLAIN QUERY PLAN 2637 SELECT plan_id FROM radroots_runtime_authored_delivery_plans 2638 WHERE state = 'retryable' AND retry_not_before_unix_ms <= 100 2639 ORDER BY state, retry_not_before_unix_ms, claim_expires_at_unix_ms, 2640 updated_at_unix_ms, plan_id", 2641 ) 2642 .fetch_all(&store.pool) 2643 .await 2644 .expect("query plan") 2645 .iter() 2646 .map(|row| row.get::<String, _>("detail")) 2647 .collect::<Vec<_>>() 2648 .join(" "); 2649 assert!(plan.contains("radroots_runtime_authored_delivery_ready_idx")); 2650 } 2651 } 2652 2653 #[cfg(test)] 2654 #[cfg_attr(coverage_nightly, coverage(off))] 2655 #[path = "authored_draft_submission_tests.rs"] 2656 mod draft_submission_tests; 2657 2658 #[cfg(test)] 2659 #[cfg_attr(coverage_nightly, coverage(off))] 2660 #[path = "authored_signed_durability_tests.rs"] 2661 mod signed_durability_tests; 2662 2663 #[cfg(test)] 2664 #[cfg_attr(coverage_nightly, coverage(off))] 2665 #[path = "authored_signed_fact_tests.rs"] 2666 mod signed_fact_tests; 2667 2668 #[cfg(test)] 2669 #[cfg_attr(coverage_nightly, coverage(off))] 2670 #[path = "authored_delivery_fact_tests.rs"] 2671 mod delivery_fact_tests; 2672 2673 #[cfg(test)] 2674 #[cfg_attr(coverage_nightly, coverage(off))] 2675 #[path = "authored_delivery_reconciliation_tests.rs"] 2676 mod delivery_reconciliation_tests; 2677 2678 #[cfg(test)] 2679 #[cfg_attr(coverage_nightly, coverage(off))] 2680 #[path = "authored_signed_fact_fixture.rs"] 2681 pub(crate) mod signed_fact_fixture;