authored_v10.rs (53811B)
1 use crate::{Error, authored, journal, outbox}; 2 use core::num::{NonZeroU32, NonZeroU64}; 3 use radroots_event::SignedEvent; 4 use radroots_event_codec::{Codec, verify}; 5 use radroots_storage::{ 6 authored::{ 7 AdmissionState, ArtifactOrigin, AuthoredArtifact, AuthoredArtifactId, AuthoredOperation, 8 FailureClass, RetrySchedule, WorkClaim, WorkFailure, WorkPhase, 9 }, 10 authored_delivery::{AuthoredDeliveryPlan, AuthoredDeliveryPlanId, AuthoredDeliveryState}, 11 journal::{JournalState, OperationRecord}, 12 outbox::{OutboxRecord, OutboxStage}, 13 }; 14 use radroots_transport::{ 15 DeliveryReceipt, 16 policy::{SatisfactionState, evaluate_satisfaction}, 17 sink::DeliveryTargetReceipt, 18 }; 19 use sha2::{Digest, Sha256}; 20 use sqlx::{Row, Sqlite, SqliteConnection}; 21 use std::collections::BTreeSet; 22 23 const BLOCKED_IDENTIFIERS_MAX: usize = 32; 24 25 #[derive(Clone, Debug, Eq, PartialEq)] 26 pub struct AuthoredV10Preflight { 27 operation_count: u64, 28 event_count: u64, 29 outbox_count: u64, 30 target_count: u64, 31 attempt_count: u64, 32 importable_count: u64, 33 prepared_or_recoverable: u64, 34 signed_without_complete_event: u64, 35 invalid_or_unsupported: u64, 36 blocked_operation_ids: Vec<[u8; 16]>, 37 source_digest: [u8; 32], 38 } 39 40 impl AuthoredV10Preflight { 41 pub const fn operation_count(&self) -> u64 { 42 self.operation_count 43 } 44 pub const fn event_count(&self) -> u64 { 45 self.event_count 46 } 47 pub const fn outbox_count(&self) -> u64 { 48 self.outbox_count 49 } 50 pub const fn target_count(&self) -> u64 { 51 self.target_count 52 } 53 pub const fn attempt_count(&self) -> u64 { 54 self.attempt_count 55 } 56 pub const fn importable_count(&self) -> u64 { 57 self.importable_count 58 } 59 pub const fn prepared_or_recoverable(&self) -> u64 { 60 self.prepared_or_recoverable 61 } 62 pub const fn signed_without_complete_event(&self) -> u64 { 63 self.signed_without_complete_event 64 } 65 pub const fn invalid_or_unsupported(&self) -> u64 { 66 self.invalid_or_unsupported 67 } 68 pub fn blocked_operation_ids(&self) -> &[[u8; 16]] { 69 self.blocked_operation_ids.as_slice() 70 } 71 pub const fn source_digest(&self) -> &[u8; 32] { 72 &self.source_digest 73 } 74 pub const fn is_eligible(&self) -> bool { 75 self.prepared_or_recoverable == 0 76 && self.signed_without_complete_event == 0 77 && self.invalid_or_unsupported == 0 78 } 79 80 pub(crate) fn blocked_error(&self) -> Error { 81 Error::AuthoredMigrationBlocked { 82 prepared_or_recoverable: self.prepared_or_recoverable, 83 signed_without_complete_event: self.signed_without_complete_event, 84 invalid_or_unsupported: self.invalid_or_unsupported, 85 } 86 } 87 } 88 89 struct Candidate { 90 journal: OperationRecord, 91 outbox: OutboxRecord, 92 event_admitted_at_unix_ms: u64, 93 } 94 95 struct EventMetadata { 96 event_id: [u8; 32], 97 contract_id: &'static str, 98 registry_version: u32, 99 } 100 101 pub(crate) struct InspectedV10 { 102 pub(crate) report: AuthoredV10Preflight, 103 candidates: Vec<Candidate>, 104 event_metadata: Vec<EventMetadata>, 105 } 106 107 pub(crate) async fn inspect(connection: &mut SqliteConnection) -> Result<InspectedV10, Error> { 108 let operation_count = count(connection, "radroots_runtime_journal_operations").await?; 109 let event_count = count(connection, "radroots_runtime_events").await?; 110 let outbox_count = count(connection, "radroots_runtime_outbox_items").await?; 111 let target_count = count(connection, "radroots_runtime_outbox_targets").await?; 112 let attempt_count = count(connection, "radroots_runtime_delivery_evidence").await?; 113 114 let mut invalid_or_unsupported = orphan_count(connection).await?; 115 let mut source_hasher = Sha256::new(); 116 source_hasher.update(b"radroots.authored.v10.preflight.v1"); 117 let event_rows = 118 sqlx::query("SELECT event_id, signed_event FROM radroots_runtime_events ORDER BY event_id") 119 .fetch_all(&mut *connection) 120 .await 121 .map_err(|source| crate::backend::startup_error(&source, metadata_error()))?; 122 let mut event_metadata = Vec::with_capacity(event_rows.len()); 123 let mut invalid_event_ids = BTreeSet::new(); 124 for row in event_rows { 125 let event_id = array::<32>(row.try_get("event_id").map_err(|_| metadata_error())?)?; 126 let raw_bytes = row 127 .try_get::<Vec<u8>, _>("signed_event") 128 .map_err(|_| metadata_error())?; 129 source_hasher.update(event_id); 130 source_hasher.update((raw_bytes.len() as u64).to_be_bytes()); 131 source_hasher.update(raw_bytes.as_slice()); 132 let Ok(raw) = String::from_utf8(raw_bytes) else { 133 invalid_event_ids.insert(event_id); 134 continue; 135 }; 136 let Ok(event) = Codec::decode_signed_event(raw.as_str()) else { 137 invalid_event_ids.insert(event_id); 138 continue; 139 }; 140 let Ok(contract) = 141 radroots_event::contract::registry_v7::validate_event_contract_registry_v7( 142 event.envelope(), 143 ) 144 else { 145 invalid_event_ids.insert(event_id); 146 continue; 147 }; 148 if event.id().as_bytes() != &event_id { 149 invalid_event_ids.insert(event_id); 150 continue; 151 } 152 event_metadata.push(EventMetadata { 153 event_id, 154 contract_id: contract.id, 155 registry_version: radroots_event::contract::RegistryVersion::CURRENT.get(), 156 }); 157 } 158 159 let mut outboxes = Vec::new(); 160 let outbox_ids = sqlx::query_scalar::<_, Vec<u8>>( 161 "SELECT item_id FROM radroots_runtime_outbox_items ORDER BY item_id", 162 ) 163 .fetch_all(&mut *connection) 164 .await 165 .map_err(|source| crate::backend::startup_error(&source, metadata_error()))?; 166 for bytes in outbox_ids { 167 source_hasher.update((bytes.len() as u64).to_be_bytes()); 168 source_hasher.update(bytes.as_slice()); 169 let item_id = match radroots_storage::outbox::OutboxItemId::new(array(bytes)?) { 170 Ok(value) => value, 171 Err(_) => { 172 invalid_or_unsupported = invalid_or_unsupported.saturating_add(1); 173 continue; 174 } 175 }; 176 match outbox::load_record(connection, item_id).await { 177 Ok(Some(record)) => outboxes.push(record), 178 Err(radroots_storage::Error::SpaceInsufficient) => { 179 return Err(Error::SpaceInsufficient); 180 } 181 Ok(None) | Err(_) => { 182 invalid_or_unsupported = invalid_or_unsupported.saturating_add(1); 183 } 184 } 185 } 186 187 let journal_rows = 188 sqlx::query("SELECT * FROM radroots_runtime_journal_operations ORDER BY instance_id") 189 .fetch_all(&mut *connection) 190 .await 191 .map_err(|source| crate::backend::startup_error(&source, metadata_error()))?; 192 let mut candidates = Vec::new(); 193 let mut matched_outboxes = BTreeSet::new(); 194 let mut prepared_or_recoverable = 0_u64; 195 let mut signed_without_complete_event = 0_u64; 196 let mut blocked_operation_ids = Vec::new(); 197 let mut referenced_event_ids = BTreeSet::new(); 198 199 for row in journal_rows { 200 let raw_id = row 201 .try_get::<Vec<u8>, _>("instance_id") 202 .map_err(|_| metadata_error())?; 203 source_hasher.update((raw_id.len() as u64).to_be_bytes()); 204 source_hasher.update(raw_id.as_slice()); 205 let Ok(record) = journal::decode_record(&row) else { 206 invalid_or_unsupported = invalid_or_unsupported.saturating_add(1); 207 push_blocked(&mut blocked_operation_ids, raw_id.as_slice()); 208 continue; 209 }; 210 let operation_id = *record.instance_id().as_bytes(); 211 match record.state() { 212 JournalState::Prepared | JournalState::Recoverable(_) => { 213 prepared_or_recoverable = prepared_or_recoverable.saturating_add(1); 214 push_blocked(&mut blocked_operation_ids, &operation_id); 215 } 216 JournalState::Signed { .. } => { 217 signed_without_complete_event = signed_without_complete_event.saturating_add(1); 218 push_blocked(&mut blocked_operation_ids, &operation_id); 219 } 220 JournalState::Committed { event_id, .. } => { 221 referenced_event_ids.insert(*event_id.as_bytes()); 222 let Some(record_outbox) = outboxes 223 .iter() 224 .find(|outbox| outbox.operation_instance_id() == record.instance_id()) 225 .cloned() 226 else { 227 signed_without_complete_event = signed_without_complete_event.saturating_add(1); 228 push_blocked(&mut blocked_operation_ids, &operation_id); 229 continue; 230 }; 231 matched_outboxes.insert(*record_outbox.item_id().as_bytes()); 232 if let Err(error) = validate_candidate(connection, &record_outbox, event_id).await { 233 if matches!(error, Error::SpaceInsufficient) { 234 return Err(error); 235 } 236 invalid_or_unsupported = invalid_or_unsupported.saturating_add(1); 237 push_blocked(&mut blocked_operation_ids, &operation_id); 238 continue; 239 } 240 let admitted_at = event_admitted_at(connection, event_id.as_bytes()).await?; 241 source_hasher.update(record.input_digest().as_bytes()); 242 source_hasher.update(record_outbox.item_id().as_bytes()); 243 source_hasher.update( 244 record_outbox 245 .request() 246 .payload() 247 .event() 248 .raw_json() 249 .as_bytes(), 250 ); 251 candidates.push(Candidate { 252 journal: record, 253 outbox: record_outbox, 254 event_admitted_at_unix_ms: admitted_at, 255 }); 256 } 257 } 258 } 259 let unmatched = outboxes 260 .iter() 261 .filter(|record| !matched_outboxes.contains(record.item_id().as_bytes())) 262 .count(); 263 let unreferenced_invalid_events = invalid_event_ids.difference(&referenced_event_ids).count(); 264 invalid_or_unsupported = invalid_or_unsupported 265 .saturating_add(u64::try_from(unmatched).unwrap_or(u64::MAX)) 266 .saturating_add(u64::try_from(unreferenced_invalid_events).unwrap_or(u64::MAX)); 267 268 let report = AuthoredV10Preflight { 269 operation_count, 270 event_count, 271 outbox_count, 272 target_count, 273 attempt_count, 274 importable_count: u64::try_from(candidates.len()).unwrap_or(u64::MAX), 275 prepared_or_recoverable, 276 signed_without_complete_event, 277 invalid_or_unsupported, 278 blocked_operation_ids, 279 source_digest: source_hasher.finalize().into(), 280 }; 281 Ok(InspectedV10 { 282 report, 283 candidates, 284 event_metadata, 285 }) 286 } 287 288 pub(crate) async fn apply( 289 transaction: &mut sqlx::Transaction<'_, Sqlite>, 290 inspected: &InspectedV10, 291 ) -> Result<(), Error> { 292 if !inspected.report.is_eligible() { 293 return Err(inspected.report.blocked_error()); 294 } 295 for metadata in &inspected.event_metadata { 296 let result = sqlx::query( 297 "UPDATE radroots_runtime_events 298 SET admitted_contract_id = ?, admitted_registry_version = ? 299 WHERE event_id = ? AND admitted_contract_id IS NULL 300 AND admitted_registry_version IS NULL", 301 ) 302 .bind(metadata.contract_id) 303 .bind(i64::from(metadata.registry_version)) 304 .bind(metadata.event_id.as_slice()) 305 .execute(&mut **transaction) 306 .await 307 .map_err(|source| crate::backend::startup_error(&source, metadata_error()))?; 308 if result.rows_affected() != 1 { 309 return Err(metadata_error()); 310 } 311 } 312 for candidate in &inspected.candidates { 313 let (operation, artifact, plan) = convert_candidate(candidate)?; 314 authored::persist_operation(transaction, &operation) 315 .await 316 .map_err(|source| match source { 317 radroots_storage::Error::SpaceInsufficient => Error::SpaceInsufficient, 318 _ => metadata_error(), 319 })?; 320 persist_imported_artifact(transaction, &artifact).await?; 321 authored::persist_plan_v11(transaction, &plan) 322 .await 323 .map_err(|source| match source { 324 radroots_storage::Error::SpaceInsufficient => Error::SpaceInsufficient, 325 _ => metadata_error(), 326 })?; 327 } 328 let operation_count = count_transaction( 329 transaction, 330 "SELECT COUNT(*) FROM radroots_runtime_authored_operations", 331 ) 332 .await?; 333 let artifact_count = count_transaction( 334 transaction, 335 "SELECT COUNT(*) FROM radroots_runtime_authored_artifacts", 336 ) 337 .await?; 338 let plan_count = count_transaction( 339 transaction, 340 "SELECT COUNT(*) FROM radroots_runtime_authored_delivery_plans", 341 ) 342 .await?; 343 if [operation_count, artifact_count, plan_count] 344 .iter() 345 .any(|count| *count != inspected.report.importable_count) 346 { 347 return Err(metadata_error()); 348 } 349 let foreign_keys = sqlx::query("PRAGMA foreign_key_check") 350 .fetch_all(&mut **transaction) 351 .await 352 .map_err(|source| crate::backend::startup_error(&source, metadata_error()))?; 353 if !foreign_keys.is_empty() { 354 return Err(metadata_error()); 355 } 356 let integrity = sqlx::query_scalar::<_, String>("PRAGMA integrity_check") 357 .fetch_all(&mut **transaction) 358 .await 359 .map_err(|source| crate::backend::startup_error(&source, metadata_error()))?; 360 if integrity.as_slice() != ["ok"] { 361 return Err(metadata_error()); 362 } 363 sqlx::query( 364 "INSERT INTO radroots_runtime_authored_migration_evidence ( 365 source_version, operation_count, event_count, outbox_count, 366 target_count, attempt_count, imported_count, source_digest, 367 completed_at_unix_ms 368 ) VALUES (10, ?, ?, ?, ?, ?, ?, ?, ?)", 369 ) 370 .bind(i64_from_u64(inspected.report.operation_count)?) 371 .bind(i64_from_u64(inspected.report.event_count)?) 372 .bind(i64_from_u64(inspected.report.outbox_count)?) 373 .bind(i64_from_u64(inspected.report.target_count)?) 374 .bind(i64_from_u64(inspected.report.attempt_count)?) 375 .bind(i64_from_u64(inspected.report.importable_count)?) 376 .bind(inspected.report.source_digest.as_slice()) 377 .bind(migration_timestamp(inspected)) 378 .execute(&mut **transaction) 379 .await 380 .map_err(|source| crate::backend::startup_error(&source, metadata_error()))?; 381 Ok(()) 382 } 383 384 fn convert_candidate( 385 candidate: &Candidate, 386 ) -> Result<(AuthoredOperation, AuthoredArtifact, AuthoredDeliveryPlan), Error> { 387 let operation_id = candidate.journal.instance_id(); 388 let artifact_id = AuthoredArtifactId::new(derive_id( 389 b"radroots.authored.v10.artifact.v1", 390 operation_id.as_bytes(), 391 )) 392 .map_err(|_| metadata_error())?; 393 let plan_id = AuthoredDeliveryPlanId::new(*candidate.outbox.item_id().as_bytes()) 394 .map_err(|_| metadata_error())?; 395 let operation = AuthoredOperation::new( 396 operation_id, 397 vec![artifact_id], 398 candidate.journal.prepared_at_unix_ms(), 399 ) 400 .map_err(|_| metadata_error())?; 401 let mut artifact = AuthoredArtifact::imported_signed( 402 artifact_id, 403 operation_id, 404 0, 405 candidate.outbox.request().payload().event().clone(), 406 candidate.journal.prepared_at_unix_ms(), 407 ) 408 .map_err(|_| metadata_error())?; 409 let admitted_at = candidate 410 .event_admitted_at_unix_ms 411 .max(candidate.journal.prepared_at_unix_ms()); 412 let admission_claim = WorkClaim::new( 413 derive_id( 414 b"radroots.authored.v10.admission-claim.v1", 415 artifact_id.as_bytes(), 416 ), 417 "migration-v10", 418 NonZeroU64::MIN, 419 admitted_at, 420 admitted_at.checked_add(1).ok_or_else(metadata_error)?, 421 artifact.revision(), 422 ) 423 .map_err(|_| metadata_error())?; 424 artifact 425 .set_admission_claim(admission_claim, admitted_at) 426 .map_err(|_| metadata_error())?; 427 artifact 428 .record_admission(AdmissionState::Inserted, None, None, admitted_at) 429 .map_err(|_| metadata_error())?; 430 431 let mut plan = AuthoredDeliveryPlan::new_bound( 432 plan_id, 433 artifact_id, 434 candidate.outbox.request().clone(), 435 candidate.outbox.created_at_unix_ms(), 436 ) 437 .map_err(|_| metadata_error())?; 438 replay_evidence(&mut plan, &candidate.outbox)?; 439 Ok((operation, artifact, plan)) 440 } 441 442 // This conversion runs immediately after v11 DDL, before successor columns exist. 443 async fn persist_imported_artifact( 444 transaction: &mut sqlx::Transaction<'_, Sqlite>, 445 artifact: &AuthoredArtifact, 446 ) -> Result<(), Error> { 447 artifact.validate().map_err(|_| metadata_error())?; 448 if artifact.origin() != ArtifactOrigin::ImportedSigned 449 || artifact.admission_state() != AdmissionState::Inserted 450 { 451 return Err(metadata_error()); 452 } 453 let signed = artifact.signed().ok_or_else(metadata_error)?; 454 sqlx::query( 455 "INSERT INTO radroots_runtime_authored_artifacts ( 456 artifact_id, operation_id, ordinal, origin, signing_state, admission_state, 457 signed_raw_json, signed_raw_sha256, 458 created_at_unix_ms, updated_at_unix_ms, revision, snapshot 459 ) VALUES (?, ?, ?, 'imported_signed', 'signed', 'inserted', ?, ?, ?, ?, ?, ?)", 460 ) 461 .bind(artifact.artifact_id().as_bytes().as_slice()) 462 .bind(artifact.operation_id().as_bytes().as_slice()) 463 .bind(i64::from(artifact.ordinal())) 464 .bind(signed.event().raw_json().as_bytes()) 465 .bind(signed.raw_json_sha256().as_slice()) 466 .bind(i64_from_u64(artifact.created_at_unix_ms())?) 467 .bind(i64_from_u64(artifact.updated_at_unix_ms())?) 468 .bind(i64_from_u64(artifact.revision().get())?) 469 .bind(authored::encode_snapshot(artifact).map_err(|_| metadata_error())?) 470 .execute(&mut **transaction) 471 .await 472 .map_err(|source| crate::backend::startup_error(&source, metadata_error()))?; 473 Ok(()) 474 } 475 476 fn replay_evidence(plan: &mut AuthoredDeliveryPlan, legacy: &OutboxRecord) -> Result<(), Error> { 477 let Some(last_attempt) = legacy.last_attempt() else { 478 if !legacy.evidence().is_empty() 479 || !matches!(legacy.stage(), OutboxStage::Pending | OutboxStage::Leased) 480 { 481 return Err(metadata_error()); 482 } 483 return Ok(()); 484 }; 485 let mut cumulative = Vec::new(); 486 for attempt_number in 1..=last_attempt.get() { 487 let evidence = legacy 488 .evidence() 489 .iter() 490 .filter(|entry| entry.attempt().get() == attempt_number) 491 .collect::<Vec<_>>(); 492 if evidence.len() != legacy.request().target_set().len() { 493 return Err(metadata_error()); 494 } 495 let recorded_at = evidence 496 .first() 497 .map(|entry| entry.recorded_at_unix_ms()) 498 .ok_or_else(metadata_error)?; 499 if evidence 500 .iter() 501 .any(|entry| entry.recorded_at_unix_ms() != recorded_at) 502 { 503 return Err(metadata_error()); 504 } 505 let mut receipts = Vec::with_capacity(evidence.len()); 506 for target in legacy.request().target_set().targets() { 507 let entry = evidence 508 .iter() 509 .find(|entry| entry.target() == target.fingerprint()) 510 .ok_or_else(metadata_error)?; 511 let receipt = if entry.was_attempted() { 512 DeliveryTargetReceipt::attempted(target.clone(), entry.outcome().clone()) 513 } else { 514 DeliveryTargetReceipt::skipped(target.clone(), entry.outcome().clone()) 515 .map_err(|_| metadata_error())? 516 }; 517 receipts.push(receipt); 518 } 519 cumulative.extend(receipts.iter().cloned()); 520 let receipt = DeliveryReceipt::for_request(legacy.request(), receipts) 521 .map_err(|_| metadata_error())?; 522 let satisfaction = evaluate_satisfaction( 523 legacy.request().satisfaction(), 524 legacy.request().target_set(), 525 cumulative 526 .iter() 527 .map(|entry| (entry.target().fingerprint(), entry.outcome())), 528 ) 529 .map_err(|_| metadata_error())?; 530 let retry = if satisfaction == SatisfactionState::Pending { 531 let next_at = legacy 532 .evidence() 533 .iter() 534 .find(|entry| entry.attempt().get() == attempt_number.saturating_add(1)) 535 .map(|entry| entry.recorded_at_unix_ms()) 536 .or(legacy.retry_not_before_unix_ms()) 537 .map_or_else( 538 || { 539 legacy 540 .updated_at_unix_ms() 541 .max(recorded_at) 542 .checked_add(1) 543 .ok_or_else(metadata_error) 544 }, 545 Ok, 546 )?; 547 if next_at <= recorded_at { 548 return Err(metadata_error()); 549 } 550 let failure = WorkFailure::new( 551 "migrated_delivery_pending", 552 WorkPhase::Delivery, 553 FailureClass::Retryable, 554 Some(next_at), 555 None, 556 ) 557 .map_err(|_| metadata_error())?; 558 Some( 559 RetrySchedule::new( 560 NonZeroU32::new(attempt_number).ok_or_else(metadata_error)?, 561 next_at, 562 failure, 563 ) 564 .map_err(|_| metadata_error())?, 565 ) 566 } else { 567 None 568 }; 569 let claim = WorkClaim::new( 570 derive_attempt_token(plan.plan_id(), attempt_number), 571 "migration-v10", 572 NonZeroU64::new(u64::from(attempt_number)).ok_or_else(metadata_error)?, 573 recorded_at, 574 recorded_at.checked_add(1).ok_or_else(metadata_error)?, 575 plan.revision(), 576 ) 577 .map_err(|_| metadata_error())?; 578 let token = *claim.token(); 579 let generation = claim.generation(); 580 let revision = claim.row_revision(); 581 plan.claim(claim, recorded_at) 582 .map_err(|_| metadata_error())?; 583 plan.apply_receipt(&token, generation, revision, receipt, retry, recorded_at) 584 .map_err(|_| metadata_error())?; 585 } 586 let state_matches = matches!( 587 (legacy.stage(), plan.state()), 588 (OutboxStage::Satisfied, AuthoredDeliveryState::Satisfied) 589 | (OutboxStage::Exhausted, AuthoredDeliveryState::Exhausted) 590 | (OutboxStage::Retryable, AuthoredDeliveryState::Retryable) 591 | (OutboxStage::Leased, AuthoredDeliveryState::Retryable) 592 ); 593 if !state_matches { 594 return Err(metadata_error()); 595 } 596 Ok(()) 597 } 598 599 async fn validate_candidate( 600 connection: &mut SqliteConnection, 601 outbox: &OutboxRecord, 602 event_id: &radroots_event::EventId, 603 ) -> Result<(), Error> { 604 let event = outbox.request().payload().event(); 605 if event.id() != event_id { 606 return Err(metadata_error()); 607 } 608 let stored = sqlx::query("SELECT signed_event FROM radroots_runtime_events WHERE event_id = ?") 609 .bind(event_id.as_bytes().as_slice()) 610 .fetch_optional(&mut *connection) 611 .await 612 .map_err(|source| crate::backend::startup_error(&source, metadata_error()))? 613 .ok_or_else(metadata_error)?; 614 if stored 615 .try_get::<Vec<u8>, _>("signed_event") 616 .map_err(|_| metadata_error())? 617 .as_slice() 618 != event.raw_json().as_bytes() 619 { 620 return Err(metadata_error()); 621 } 622 verify_exact_event(event)?; 623 Ok(()) 624 } 625 626 fn verify_exact_event(event: &SignedEvent) -> Result<(), Error> { 627 let raw = verify::RawEvent::new(event.envelope().clone()); 628 let id = verify::id(raw).map_err(|_| metadata_error())?; 629 let signature = 630 verify::signature(id, &verify::Nip01SignatureVerifier).map_err(|_| metadata_error())?; 631 verify::contract(signature).map_err(|_| metadata_error())?; 632 Ok(()) 633 } 634 635 async fn event_admitted_at( 636 connection: &mut SqliteConnection, 637 event_id: &[u8; 32], 638 ) -> Result<u64, Error> { 639 let value = sqlx::query_scalar::<_, i64>( 640 "SELECT admitted_at_unix_ms FROM radroots_runtime_events WHERE event_id = ?", 641 ) 642 .bind(event_id.as_slice()) 643 .fetch_one(&mut *connection) 644 .await 645 .map_err(|source| crate::backend::startup_error(&source, metadata_error()))?; 646 u64_from_i64(value) 647 } 648 649 async fn orphan_count(connection: &mut SqliteConnection) -> Result<u64, Error> { 650 let value = sqlx::query_scalar::<_, i64>( 651 "SELECT 652 (SELECT COUNT(*) FROM radroots_runtime_outbox_targets AS target 653 LEFT JOIN radroots_runtime_outbox_items AS item ON item.item_id = target.item_id 654 WHERE item.item_id IS NULL) 655 + 656 (SELECT COUNT(*) FROM radroots_runtime_delivery_evidence AS evidence 657 LEFT JOIN radroots_runtime_outbox_targets AS target 658 ON target.item_id = evidence.item_id 659 AND target.target_fingerprint = evidence.target_fingerprint 660 WHERE target.item_id IS NULL)", 661 ) 662 .fetch_one(&mut *connection) 663 .await 664 .map_err(|source| crate::backend::startup_error(&source, metadata_error()))?; 665 u64_from_i64(value) 666 } 667 668 async fn count(connection: &mut SqliteConnection, table: &'static str) -> Result<u64, Error> { 669 let query = match table { 670 "radroots_runtime_journal_operations" => { 671 "SELECT COUNT(*) FROM radroots_runtime_journal_operations" 672 } 673 "radroots_runtime_events" => "SELECT COUNT(*) FROM radroots_runtime_events", 674 "radroots_runtime_outbox_items" => "SELECT COUNT(*) FROM radroots_runtime_outbox_items", 675 "radroots_runtime_outbox_targets" => "SELECT COUNT(*) FROM radroots_runtime_outbox_targets", 676 "radroots_runtime_delivery_evidence" => { 677 "SELECT COUNT(*) FROM radroots_runtime_delivery_evidence" 678 } 679 _ => return Err(metadata_error()), 680 }; 681 u64_from_i64( 682 sqlx::query_scalar::<_, i64>(query) 683 .fetch_one(&mut *connection) 684 .await 685 .map_err(|source| crate::backend::startup_error(&source, metadata_error()))?, 686 ) 687 } 688 689 async fn count_transaction( 690 transaction: &mut sqlx::Transaction<'_, Sqlite>, 691 query: &'static str, 692 ) -> Result<u64, Error> { 693 u64_from_i64( 694 sqlx::query_scalar::<_, i64>(query) 695 .fetch_one(&mut **transaction) 696 .await 697 .map_err(|source| crate::backend::startup_error(&source, metadata_error()))?, 698 ) 699 } 700 701 fn migration_timestamp(inspected: &InspectedV10) -> i64 { 702 inspected 703 .candidates 704 .iter() 705 .map(|candidate| candidate.outbox.updated_at_unix_ms()) 706 .max() 707 .unwrap_or(1) 708 .try_into() 709 .unwrap_or(i64::MAX) 710 } 711 712 fn derive_attempt_token(plan_id: AuthoredDeliveryPlanId, attempt: u32) -> [u8; 16] { 713 let mut hasher = Sha256::new(); 714 hasher.update(b"radroots.authored.v10.delivery-claim.v1"); 715 hasher.update(plan_id.as_bytes()); 716 hasher.update(attempt.to_be_bytes()); 717 let digest: [u8; 32] = hasher.finalize().into(); 718 let mut result = [0; 16]; 719 result.copy_from_slice(&digest[..16]); 720 result 721 } 722 723 fn derive_id(domain: &[u8], input: &[u8; 16]) -> [u8; 16] { 724 let mut hasher = Sha256::new(); 725 hasher.update(domain); 726 hasher.update(input); 727 let digest: [u8; 32] = hasher.finalize().into(); 728 let mut result = [0; 16]; 729 result.copy_from_slice(&digest[..16]); 730 result 731 } 732 733 fn push_blocked(target: &mut Vec<[u8; 16]>, bytes: &[u8]) { 734 if target.len() < BLOCKED_IDENTIFIERS_MAX 735 && let Ok(id) = <[u8; 16]>::try_from(bytes) 736 { 737 target.push(id); 738 } 739 } 740 741 fn array<const N: usize>(bytes: Vec<u8>) -> Result<[u8; N], Error> { 742 bytes.try_into().map_err(|_| metadata_error()) 743 } 744 745 fn i64_from_u64(value: u64) -> Result<i64, Error> { 746 i64::try_from(value).map_err(|_| metadata_error()) 747 } 748 749 fn u64_from_i64(value: i64) -> Result<u64, Error> { 750 u64::try_from(value).map_err(|_| metadata_error()) 751 } 752 753 const fn metadata_error() -> Error { 754 Error::SchemaMigrationFailed { 755 database: "runtime.sqlite", 756 target_version: 11, 757 } 758 } 759 760 #[cfg(test)] 761 #[cfg_attr(coverage_nightly, coverage(off))] 762 mod tests { 763 use super::*; 764 use crate::{SqliteStorage, migration::runtime}; 765 use radroots_storage::{ 766 Journal, Outbox, 767 event::SourceGeneration, 768 journal::{ 769 IdempotencyDigest, IdempotencyKey, JournalTransition, OperationId, OperationInstanceId, 770 PrepareOperation, 771 }, 772 outbox::{ 773 ClaimOutboxItems, DeliveryAttempt, DeliveryAttemptEvidence, DeliveryOutcome, 774 DeliveryPlanDigest, EnqueueOutboxItem, LeaseId, LeaseOwner, OutboxItemId, 775 }, 776 status::EventStoreMode, 777 }; 778 use radroots_transport::{ 779 DeliveryRequest, Target, TargetSet, 780 policy::{SatisfactionClass, SatisfactionPolicy, TargetPolicy}, 781 sink::DeliveryPayload, 782 }; 783 use sqlx::{Connection, sqlite::SqlitePoolOptions}; 784 785 const VALID_EVENT: &str = r#"{"id":"762bee187e9e645b81ec26ade05a69b5e8398caf527be8de0d9a45311ed0c7a0","pubkey":"585591529da0bab31b3b1b1f986611cf5f435dca84f978c89ee8a40cca7103df","created_at":1800000100,"kind":0,"tags":[],"content":"{\"display_name\":\"Moss Street Farm\",\"bot\":false,\"website\":\"https://mossstreet.example\",\"picture\":42}","sig":"4290da0bb6422986647bc8cd5f63bd52d49f41e7b665d3b47105b8109183e8d596f322c531d4061df53e1d2b70fda12d5d1c14f3720d7a56d9d0a03746af5109"}"#; 786 787 async fn v10_store() -> SqliteStorage { 788 let generation = SourceGeneration::new([61; 32]).expect("generation"); 789 let pool = SqlitePoolOptions::new() 790 .max_connections(1) 791 .connect("sqlite::memory:") 792 .await 793 .expect("memory SQLite"); 794 sqlx::query("PRAGMA foreign_keys = ON") 795 .execute(&pool) 796 .await 797 .expect("foreign keys"); 798 for migration in runtime::MIGRATIONS.iter().take(10) { 799 sqlx::raw_sql(runtime::migration_sql(migration.version()).expect("migration SQL")) 800 .execute(&pool) 801 .await 802 .expect("runtime migration"); 803 } 804 sqlx::raw_sql("PRAGMA application_id = 1380209236; PRAGMA user_version = 10") 805 .execute(&pool) 806 .await 807 .expect("schema metadata"); 808 sqlx::query( 809 "INSERT INTO radroots_runtime_source_generations ( 810 generation, sequence_head, state, created_at_unix_ms, retired_at_unix_ms 811 ) VALUES (?, 0, 'active', 1, NULL)", 812 ) 813 .bind(generation.as_bytes().as_slice()) 814 .execute(&pool) 815 .await 816 .expect("source generation"); 817 SqliteStorage::new(pool, generation, EventStoreMode::ReadWrite) 818 } 819 820 fn event() -> SignedEvent { 821 Codec::decode_signed_event(VALID_EVENT).expect("valid signed event") 822 } 823 824 fn operation_id(byte: u8) -> OperationInstanceId { 825 OperationInstanceId::new([byte; 16]).expect("operation") 826 } 827 828 async fn prepare(store: &SqliteStorage, byte: u8) -> OperationRecord { 829 store 830 .prepare( 831 PrepareOperation::new( 832 operation_id(byte), 833 OperationId::SyncPush, 834 IdempotencyKey::parse(format!("authored-v10-{byte}")).expect("idempotency key"), 835 IdempotencyDigest::new([byte; 32]), 836 100, 837 ) 838 .expect("prepare operation"), 839 ) 840 .await 841 .expect("prepare") 842 .record() 843 .clone() 844 } 845 846 fn request(event: SignedEvent) -> DeliveryRequest { 847 DeliveryRequest::new( 848 "migration-delivery", 849 DeliveryPayload::new(event), 850 TargetSet::new(vec![ 851 Target::nostr_relay("wss://migration.example").expect("target"), 852 ]) 853 .expect("target set"), 854 SatisfactionPolicy::new(SatisfactionClass::Accepted, TargetPolicy::any()), 855 10_000, 856 ) 857 .expect("delivery request") 858 } 859 860 async fn seed_complete(store: &SqliteStorage, byte: u8, event: SignedEvent) { 861 let prepared = prepare(store, byte).await; 862 let signed = store 863 .transition(JournalTransition::signed( 864 prepared.instance_id(), 865 prepared.revision(), 866 *event.id(), 867 )) 868 .await 869 .expect("signed transition"); 870 let enqueue = EnqueueOutboxItem::new( 871 OutboxItemId::new([byte.saturating_add(20); 16]).expect("outbox item"), 872 prepared.instance_id(), 873 DeliveryPlanDigest::new([byte.saturating_add(30); 32]), 874 request(event.clone()), 875 102, 876 ) 877 .expect("enqueue"); 878 sqlx::query( 879 "UPDATE radroots_runtime_source_generations 880 SET sequence_head = sequence_head + 1 881 WHERE generation = ?", 882 ) 883 .bind(store.generation.as_bytes().as_slice()) 884 .execute(&store.pool) 885 .await 886 .expect("advance event sequence"); 887 sqlx::query( 888 "INSERT INTO radroots_runtime_events ( 889 source_generation, source_sequence, event_id, admission_stage, 890 signed_event, admitted_at_unix_ms, updated_at_unix_ms 891 ) VALUES (?, 1, ?, 'raw', ?, 102, 102)", 892 ) 893 .bind(store.generation.as_bytes().as_slice()) 894 .bind(event.id().as_bytes().as_slice()) 895 .bind(event.raw_json().as_bytes()) 896 .execute(&store.pool) 897 .await 898 .expect("legacy event admission"); 899 store.enqueue(enqueue).await.expect("outbox enqueue"); 900 store 901 .transition(JournalTransition::committed( 902 prepared.instance_id(), 903 signed.revision(), 904 *event.id(), 905 102, 906 )) 907 .await 908 .expect("committed transition"); 909 } 910 911 async fn record_legacy_attempt( 912 store: &SqliteStorage, 913 outcome: DeliveryOutcome, 914 attempted: bool, 915 ) -> OutboxRecord { 916 let claimed = store 917 .claim( 918 ClaimOutboxItems::new( 919 LeaseOwner::parse("migration-worker").expect("lease owner"), 920 LeaseId::new([91; 16]).expect("lease id"), 921 103, 922 200, 923 1, 924 ) 925 .expect("claim request"), 926 ) 927 .await 928 .expect("claim") 929 .pop() 930 .expect("claimed item"); 931 let target = claimed.record().request().target_set().targets()[0].clone(); 932 let target_receipt = if attempted { 933 DeliveryTargetReceipt::attempted(target, outcome) 934 } else { 935 DeliveryTargetReceipt::skipped(target, outcome).expect("skipped receipt") 936 }; 937 let receipt = 938 DeliveryReceipt::for_request(claimed.record().request(), vec![target_receipt]) 939 .expect("receipt"); 940 store 941 .record_attempt( 942 DeliveryAttemptEvidence::new( 943 claimed.record().item_id(), 944 claimed.lease().id(), 945 claimed.record().revision(), 946 DeliveryAttempt::FIRST, 947 receipt, 948 150, 949 ) 950 .expect("attempt evidence"), 951 ) 952 .await 953 .expect("record attempt") 954 } 955 956 #[tokio::test] 957 async fn complete_rows_preflight_and_migrate_with_exact_signed_bytes_and_evidence() { 958 let store = v10_store().await; 959 let event = event(); 960 seed_complete(&store, 1, event.clone()).await; 961 let mut connection = store.pool.acquire().await.expect("connection"); 962 let inspected = inspect(&mut connection).await.expect("preflight"); 963 assert!(inspected.report.is_eligible()); 964 assert_eq!(inspected.report.importable_count(), 1); 965 assert_eq!(inspected.report.event_count(), 1); 966 assert_eq!( 967 sqlx::query_scalar::<_, i64>("PRAGMA user_version") 968 .fetch_one(&mut *connection) 969 .await 970 .expect("version after preflight"), 971 10 972 ); 973 crate::migration::migrate_runtime(&mut connection, crate::OpenMode::ReadWriteExisting) 974 .await 975 .expect("migration"); 976 assert_eq!( 977 sqlx::query_scalar::<_, i64>("PRAGMA user_version") 978 .fetch_one(&mut *connection) 979 .await 980 .expect("version after migration"), 981 i64::from(crate::migration::runtime::CURRENT_VERSION) 982 ); 983 let raw = sqlx::query_scalar::<_, Vec<u8>>( 984 "SELECT signed_raw_json FROM radroots_runtime_authored_artifacts", 985 ) 986 .fetch_one(&mut *connection) 987 .await 988 .expect("imported signed bytes"); 989 assert_eq!(raw.as_slice(), event.raw_json().as_bytes()); 990 assert_eq!( 991 sqlx::query_scalar::<_, i64>( 992 "SELECT imported_count FROM radroots_runtime_authored_migration_evidence", 993 ) 994 .fetch_one(&mut *connection) 995 .await 996 .expect("migration evidence"), 997 1 998 ); 999 assert_eq!( 1000 sqlx::query_scalar::<_, String>( 1001 "SELECT admitted_contract_id FROM radroots_runtime_events", 1002 ) 1003 .fetch_one(&mut *connection) 1004 .await 1005 .expect("contract metadata"), 1006 "radroots.profile.metadata.v1" 1007 ); 1008 } 1009 1010 #[tokio::test] 1011 async fn prepared_and_event_id_only_rows_block_without_schema_mutation() { 1012 let store = v10_store().await; 1013 let prepared = prepare(&store, 2).await; 1014 let mut connection = store.pool.acquire().await.expect("connection"); 1015 let report = inspect(&mut connection) 1016 .await 1017 .expect("prepared preflight") 1018 .report; 1019 assert!(!report.is_eligible()); 1020 assert_eq!(report.prepared_or_recoverable(), 1); 1021 assert_eq!( 1022 report.blocked_operation_ids(), 1023 &[prepared.instance_id().as_bytes().to_owned()] 1024 ); 1025 assert!(matches!( 1026 crate::migration::migrate_runtime(&mut connection, crate::OpenMode::ReadWriteExisting) 1027 .await, 1028 Err(Error::AuthoredMigrationBlocked { 1029 prepared_or_recoverable: 1, 1030 .. 1031 }) 1032 )); 1033 assert_eq!( 1034 sqlx::query_scalar::<_, i64>("PRAGMA user_version") 1035 .fetch_one(&mut *connection) 1036 .await 1037 .expect("preserved version"), 1038 10 1039 ); 1040 assert_eq!( 1041 sqlx::query_scalar::<_, i64>( 1042 "SELECT COUNT(*) FROM sqlite_schema 1043 WHERE name = 'radroots_runtime_authored_operations'", 1044 ) 1045 .fetch_one(&mut *connection) 1046 .await 1047 .expect("no successor table"), 1048 0 1049 ); 1050 1051 drop(connection); 1052 let store = v10_store().await; 1053 let prepared = prepare(&store, 3).await; 1054 store 1055 .transition(JournalTransition::signed( 1056 prepared.instance_id(), 1057 prepared.revision(), 1058 *event().id(), 1059 )) 1060 .await 1061 .expect("signed transition"); 1062 let mut connection = store.pool.acquire().await.expect("connection"); 1063 let report = inspect(&mut connection) 1064 .await 1065 .expect("signed preflight") 1066 .report; 1067 assert!(!report.is_eligible()); 1068 assert_eq!(report.signed_without_complete_event(), 1); 1069 } 1070 1071 #[tokio::test] 1072 async fn invalid_signature_blocks_complete_rows_and_retains_v10_authority() { 1073 let store = v10_store().await; 1074 let valid = event(); 1075 let mut wire = valid.wire().clone(); 1076 wire.sig = "dd".repeat(64); 1077 let raw = serde_json::to_string(&wire).expect("invalid raw"); 1078 let invalid = SignedEvent::from_wire_verified_id(wire, raw).expect("id-valid event"); 1079 seed_complete(&store, 4, invalid).await; 1080 let mut connection = store.pool.acquire().await.expect("connection"); 1081 let report = inspect(&mut connection).await.expect("preflight").report; 1082 assert!(!report.is_eligible()); 1083 assert_eq!(report.invalid_or_unsupported(), 1); 1084 assert_eq!( 1085 sqlx::query_scalar::<_, i64>("PRAGMA user_version") 1086 .fetch_one(&mut *connection) 1087 .await 1088 .expect("preserved version"), 1089 10 1090 ); 1091 1092 drop(connection); 1093 let store = v10_store().await; 1094 let valid = event(); 1095 seed_complete(&store, 14, valid.clone()).await; 1096 let mut wire = valid.wire().clone(); 1097 wire.sig = "ee".repeat(64); 1098 let mismatched_raw = serde_json::to_vec(&wire).expect("mismatched raw"); 1099 sqlx::query("DROP TRIGGER radroots_runtime_events_raw_update_guard") 1100 .execute(&store.pool) 1101 .await 1102 .expect("remove update guard for corruption fixture"); 1103 sqlx::query("UPDATE radroots_runtime_events SET signed_event = ? WHERE event_id = ?") 1104 .bind(mismatched_raw) 1105 .bind(valid.id().as_bytes().as_slice()) 1106 .execute(&store.pool) 1107 .await 1108 .expect("stored raw mismatch fixture"); 1109 let mut connection = store.pool.acquire().await.expect("connection"); 1110 let report = inspect(&mut connection).await.expect("preflight").report; 1111 assert_eq!(report.invalid_or_unsupported(), 1); 1112 } 1113 1114 #[tokio::test] 1115 async fn malformed_event_encodings_and_identifiers_are_counted_fail_closed() { 1116 for signed_event in [vec![0xff], b"not-json".to_vec()] { 1117 let store = v10_store().await; 1118 sqlx::query( 1119 "INSERT INTO radroots_runtime_events ( 1120 source_generation, source_sequence, event_id, admission_stage, 1121 signed_event, admitted_at_unix_ms, updated_at_unix_ms 1122 ) VALUES (?, 1, ?, 'raw', ?, 10, 10)", 1123 ) 1124 .bind(store.generation.as_bytes().as_slice()) 1125 .bind([72; 32].as_slice()) 1126 .bind(signed_event) 1127 .execute(&store.pool) 1128 .await 1129 .expect("malformed event fixture"); 1130 let mut connection = store.pool.acquire().await.expect("connection"); 1131 let report = inspect(&mut connection).await.expect("preflight").report; 1132 assert_eq!(report.invalid_or_unsupported(), 1); 1133 } 1134 1135 let store = v10_store().await; 1136 sqlx::query( 1137 "INSERT INTO radroots_runtime_events ( 1138 source_generation, source_sequence, event_id, admission_stage, 1139 signed_event, admitted_at_unix_ms, updated_at_unix_ms 1140 ) VALUES (?, 1, ?, 'raw', ?, 10, 10)", 1141 ) 1142 .bind(store.generation.as_bytes().as_slice()) 1143 .bind([73; 32].as_slice()) 1144 .bind(event().raw_json().as_bytes()) 1145 .execute(&store.pool) 1146 .await 1147 .expect("mismatched event fixture"); 1148 let mut connection = store.pool.acquire().await.expect("connection"); 1149 let report = inspect(&mut connection).await.expect("preflight").report; 1150 assert_eq!(report.invalid_or_unsupported(), 1); 1151 } 1152 1153 #[tokio::test] 1154 async fn committed_without_outbox_and_corrupt_journal_rows_are_blocked() { 1155 let store = v10_store().await; 1156 let prepared = prepare(&store, 6).await; 1157 let signed = store 1158 .transition(JournalTransition::signed( 1159 prepared.instance_id(), 1160 prepared.revision(), 1161 *event().id(), 1162 )) 1163 .await 1164 .expect("signed transition"); 1165 store 1166 .transition(JournalTransition::committed( 1167 prepared.instance_id(), 1168 signed.revision(), 1169 *event().id(), 1170 102, 1171 )) 1172 .await 1173 .expect("committed transition"); 1174 let mut connection = store.pool.acquire().await.expect("connection"); 1175 let report = inspect(&mut connection).await.expect("preflight").report; 1176 assert_eq!(report.signed_without_complete_event(), 1); 1177 1178 drop(connection); 1179 let store = v10_store().await; 1180 prepare(&store, 7).await; 1181 sqlx::query( 1182 "UPDATE radroots_runtime_journal_operations SET operation_id = X'FF' 1183 WHERE instance_id = ?", 1184 ) 1185 .bind(operation_id(7).as_bytes().as_slice()) 1186 .execute(&store.pool) 1187 .await 1188 .expect("corrupt journal fixture"); 1189 let mut connection = store.pool.acquire().await.expect("connection"); 1190 let report = inspect(&mut connection).await.expect("preflight").report; 1191 assert_eq!(report.invalid_or_unsupported(), 1); 1192 } 1193 1194 #[tokio::test] 1195 async fn eligible_and_retryable_legacy_evidence_replays_exactly() { 1196 for (byte, outcome, attempted) in [ 1197 (8, DeliveryOutcome::accepted(), true), 1198 (9, DeliveryOutcome::unavailable(), true), 1199 (10, DeliveryOutcome::unavailable(), false), 1200 ] { 1201 let store = v10_store().await; 1202 seed_complete(&store, byte, event()).await; 1203 let legacy = record_legacy_attempt(&store, outcome, attempted).await; 1204 assert!(matches!( 1205 legacy.stage(), 1206 OutboxStage::Satisfied | OutboxStage::Retryable 1207 )); 1208 let mut connection = store.pool.acquire().await.expect("connection"); 1209 let inspected = inspect(&mut connection).await.expect("preflight"); 1210 assert!(inspected.report.is_eligible()); 1211 for candidate in &inspected.candidates { 1212 convert_candidate(candidate).expect("candidate conversion"); 1213 } 1214 crate::migration::migrate_runtime(&mut connection, crate::OpenMode::ReadWriteExisting) 1215 .await 1216 .expect("migration with delivery evidence"); 1217 assert_eq!( 1218 sqlx::query_scalar::<_, i64>( 1219 "SELECT attempt_count FROM radroots_runtime_authored_delivery_plans", 1220 ) 1221 .fetch_one(&mut *connection) 1222 .await 1223 .expect("attempt count"), 1224 1 1225 ); 1226 } 1227 } 1228 1229 #[tokio::test] 1230 async fn apply_rejects_ineligible_preflight_and_blocked_ids_are_bounded() { 1231 let store = v10_store().await; 1232 prepare(&store, 10).await; 1233 let mut connection = store.pool.acquire().await.expect("connection"); 1234 let inspected = inspect(&mut connection).await.expect("preflight"); 1235 let mut transaction = connection.begin().await.expect("transaction"); 1236 assert!(matches!( 1237 apply(&mut transaction, &inspected).await, 1238 Err(Error::AuthoredMigrationBlocked { .. }) 1239 )); 1240 transaction.rollback().await.expect("rollback"); 1241 1242 let mut blocked = Vec::new(); 1243 push_blocked(&mut blocked, &[0; 15]); 1244 assert!(blocked.is_empty()); 1245 for byte in 0..=BLOCKED_IDENTIFIERS_MAX { 1246 push_blocked(&mut blocked, &[u8::try_from(byte).unwrap_or(u8::MAX); 16]); 1247 } 1248 assert_eq!(blocked.len(), BLOCKED_IDENTIFIERS_MAX); 1249 } 1250 1251 #[tokio::test] 1252 async fn apply_detects_metadata_races_and_import_count_drift() { 1253 let store = v10_store().await; 1254 seed_complete(&store, 11, event()).await; 1255 let mut connection = store.pool.acquire().await.expect("connection"); 1256 let inspected = inspect(&mut connection).await.expect("preflight"); 1257 let mut transaction = connection.begin().await.expect("transaction"); 1258 sqlx::raw_sql(runtime::migration_sql(11).expect("migration SQL")) 1259 .execute(&mut *transaction) 1260 .await 1261 .expect("successor schema"); 1262 sqlx::query( 1263 "UPDATE radroots_runtime_events 1264 SET admitted_contract_id = 'radroots.profile.metadata.v1', 1265 admitted_registry_version = 7", 1266 ) 1267 .execute(&mut *transaction) 1268 .await 1269 .expect("concurrent metadata fixture"); 1270 assert!(matches!( 1271 apply(&mut transaction, &inspected).await, 1272 Err(Error::SchemaMigrationFailed { .. }) 1273 )); 1274 transaction.rollback().await.expect("rollback"); 1275 1276 drop(connection); 1277 let store = v10_store().await; 1278 seed_complete(&store, 12, event()).await; 1279 let mut connection = store.pool.acquire().await.expect("connection"); 1280 let mut inspected = inspect(&mut connection).await.expect("preflight"); 1281 inspected.candidates.clear(); 1282 let mut transaction = connection.begin().await.expect("transaction"); 1283 sqlx::raw_sql(runtime::migration_sql(11).expect("migration SQL")) 1284 .execute(&mut *transaction) 1285 .await 1286 .expect("successor schema"); 1287 assert!(matches!( 1288 apply(&mut transaction, &inspected).await, 1289 Err(Error::SchemaMigrationFailed { .. }) 1290 )); 1291 transaction.rollback().await.expect("rollback"); 1292 } 1293 1294 #[tokio::test] 1295 async fn apply_rejects_successor_foreign_key_corruption() { 1296 let store = v10_store().await; 1297 let mut connection = store.pool.acquire().await.expect("connection"); 1298 let inspected = inspect(&mut connection).await.expect("preflight"); 1299 sqlx::query("PRAGMA foreign_keys = OFF") 1300 .execute(&mut *connection) 1301 .await 1302 .expect("disable fixture foreign keys"); 1303 let mut transaction = connection.begin().await.expect("transaction"); 1304 sqlx::raw_sql(runtime::migration_sql(11).expect("migration SQL")) 1305 .execute(&mut *transaction) 1306 .await 1307 .expect("successor schema"); 1308 sqlx::query( 1309 "INSERT INTO radroots_runtime_authored_delivery_targets ( 1310 plan_id, ordinal, target_fingerprint, target_snapshot 1311 ) VALUES (?, 0, 'orphan', x'7b7d')", 1312 ) 1313 .bind([99; 16].as_slice()) 1314 .execute(&mut *transaction) 1315 .await 1316 .expect("orphan fixture"); 1317 assert!(matches!( 1318 apply(&mut transaction, &inspected).await, 1319 Err(Error::SchemaMigrationFailed { .. }) 1320 )); 1321 transaction.rollback().await.expect("rollback"); 1322 } 1323 1324 #[tokio::test] 1325 async fn nonadvancing_legacy_retry_schedule_is_rejected() { 1326 let store = v10_store().await; 1327 seed_complete(&store, 13, event()).await; 1328 let legacy = record_legacy_attempt(&store, DeliveryOutcome::unavailable(), true).await; 1329 sqlx::query( 1330 "UPDATE radroots_runtime_outbox_items SET retry_not_before_unix_ms = 150 1331 WHERE item_id = ?", 1332 ) 1333 .bind(legacy.item_id().as_bytes().as_slice()) 1334 .execute(&store.pool) 1335 .await 1336 .expect("nonadvancing retry fixture"); 1337 let mut connection = store.pool.acquire().await.expect("connection"); 1338 let inspected = inspect(&mut connection).await.expect("preflight"); 1339 assert!(inspected.report.is_eligible()); 1340 assert!(matches!( 1341 convert_candidate(&inspected.candidates[0]), 1342 Err(Error::SchemaMigrationFailed { .. }) 1343 )); 1344 } 1345 1346 #[tokio::test] 1347 async fn interrupted_conversion_rolls_back_schema_and_imported_rows() { 1348 let store = v10_store().await; 1349 seed_complete(&store, 5, event()).await; 1350 let mut connection = store.pool.acquire().await.expect("connection"); 1351 let inspected = inspect(&mut connection).await.expect("preflight"); 1352 let mut transaction = connection 1353 .begin_with("BEGIN IMMEDIATE") 1354 .await 1355 .expect("migration transaction"); 1356 sqlx::raw_sql(runtime::migration_sql(11).expect("migration SQL")) 1357 .execute(&mut *transaction) 1358 .await 1359 .expect("successor schema"); 1360 sqlx::raw_sql( 1361 "CREATE TEMP TRIGGER radroots_test_interrupt_authored_import 1362 BEFORE INSERT ON radroots_runtime_authored_artifacts 1363 BEGIN SELECT RAISE(ABORT, 'injected migration interruption'); END", 1364 ) 1365 .execute(&mut *transaction) 1366 .await 1367 .expect("fault trigger"); 1368 assert!(apply(&mut transaction, &inspected).await.is_err()); 1369 transaction.rollback().await.expect("rollback"); 1370 assert_eq!( 1371 sqlx::query_scalar::<_, i64>("PRAGMA user_version") 1372 .fetch_one(&mut *connection) 1373 .await 1374 .expect("preserved version"), 1375 10 1376 ); 1377 assert_eq!( 1378 sqlx::query_scalar::<_, i64>( 1379 "SELECT COUNT(*) FROM sqlite_schema 1380 WHERE name = 'radroots_runtime_authored_operations'", 1381 ) 1382 .fetch_one(&mut *connection) 1383 .await 1384 .expect("rolled-back schema"), 1385 0 1386 ); 1387 } 1388 }