state_recovery.rs (22305B)
1 //! Bounded restart recovery for committed exact-byte delivery work. 2 3 use core::fmt; 4 5 use radroots_service_sqlite::{ 6 ServiceSqliteTransaction, ServiceSqliteTransactionError, ServiceSqliteTransactionErrorKind, 7 }; 8 use sha2::{Digest, Sha256}; 9 use sqlx::Row; 10 11 use crate::state_delivery::{ 12 DeliveryOperationError, MycDeliveryAttemptId, MycDeliveryAttemptStatus, MycDeliveryJobId, 13 MycDeliveryJobRecord, MycDeliveryJobStatus, MycDeliveryRetryJitter, MycDeliverySourceKind, 14 MycDeliveryTargetStatus, MycDeliveryTimeUnixMs, finalize_job_if_terminal, read_attempts, 15 read_job, recover_expired, 16 }; 17 use crate::state_discovery::{ 18 MycDiscoveryPolicies, promote_current_if_desired, verify_document_for_delivery_job, 19 }; 20 use crate::state_repository::{ 21 MycStateRepository, MycStateRepositoryError, MycStateRepositoryErrorKind, PersistedMetadata, 22 RepositoryOperationError, require_expected_metadata, 23 }; 24 use crate::state_response::verify_response_for_delivery_job; 25 26 /// Maximum delivery jobs examined in one restart-recovery transaction. 27 pub const MYC_DELIVERY_RECOVERY_BATCH_MAX_COUNT: usize = 128; 28 29 const RECOVERY_JITTER_DOMAIN: &[u8] = b"radroots.myc.delivery_recovery_jitter.v1\0"; 30 31 const READ_INVARIANTS_SQL: &str = r#"SELECT 32 (SELECT COUNT(*) FROM delivery_jobs j 33 WHERE (j.source_kind = 'signer_response' AND NOT EXISTS ( 34 SELECT 1 FROM nip46_signed_responses r 35 WHERE r.operation_id = j.source_id 36 AND r.response_sha256 = j.artifact_sha256 37 ) AND NOT EXISTS ( 38 SELECT 1 FROM nip46_pending_responses r 39 WHERE r.operation_id = j.source_id 40 AND r.response_sha256 = j.artifact_sha256 41 )) 42 OR (j.source_kind = 'discovery_handler' AND NOT EXISTS ( 43 SELECT 1 FROM discovery_documents d 44 WHERE d.generation_id = j.source_id 45 AND d.event_sha256 = j.artifact_sha256 46 ))) AS invalid_sources, 47 (SELECT COUNT(*) FROM ( 48 SELECT operation_id FROM nip46_signed_responses 49 UNION ALL 50 SELECT operation_id FROM nip46_pending_responses 51 ) r 52 WHERE NOT EXISTS ( 53 SELECT 1 FROM delivery_jobs j 54 WHERE j.source_kind = 'signer_response' AND j.source_id = r.operation_id 55 )) AS orphan_responses, 56 (SELECT COUNT(*) FROM discovery_documents d 57 WHERE NOT EXISTS ( 58 SELECT 1 FROM delivery_jobs j 59 WHERE j.source_kind = 'discovery_handler' AND j.source_id = d.generation_id 60 )) AS orphan_documents, 61 (SELECT COUNT(*) FROM delivery_targets t 62 WHERE NOT EXISTS (SELECT 1 FROM delivery_jobs j WHERE j.job_id = t.job_id)) 63 AS orphan_targets, 64 (SELECT COUNT(*) FROM delivery_attempts a 65 WHERE NOT EXISTS ( 66 SELECT 1 FROM delivery_targets t 67 WHERE t.job_id = a.job_id AND t.target_index = a.target_index 68 )) AS orphan_attempts, 69 (SELECT COUNT(*) FROM delivery_jobs WHERE status IN ('pending', 'active')) AS active_jobs"#; 70 71 const READ_FIRST_JOB_IDS_SQL: &str = r#"SELECT 72 CASE WHEN typeof(job_id) = 'blob' AND length(job_id) = 32 73 THEN job_id ELSE NULL END AS job_id 74 FROM delivery_jobs 75 WHERE status IN ('pending', 'active') 76 OR (status = 'delivered' AND source_kind = 'discovery_handler' AND EXISTS ( 77 SELECT 1 FROM discovery_publication_state s 78 WHERE s.singleton = 1 79 AND s.desired_job_id = delivery_jobs.job_id 80 AND (s.current_job_id IS NULL OR s.current_job_id != s.desired_job_id) 81 )) 82 ORDER BY job_id 83 LIMIT 129"#; 84 85 const READ_NEXT_JOB_IDS_SQL: &str = r#"SELECT 86 CASE WHEN typeof(job_id) = 'blob' AND length(job_id) = 32 87 THEN job_id ELSE NULL END AS job_id 88 FROM delivery_jobs 89 WHERE (status IN ('pending', 'active') 90 OR (status = 'delivered' AND source_kind = 'discovery_handler' AND EXISTS ( 91 SELECT 1 FROM discovery_publication_state s 92 WHERE s.singleton = 1 93 AND s.desired_job_id = delivery_jobs.job_id 94 AND (s.current_job_id IS NULL OR s.current_job_id != s.desired_job_id) 95 ))) AND job_id > ? 96 ORDER BY job_id 97 LIMIT 129"#; 98 99 /// One-use entropy supplied by the runtime recovery boundary. 100 pub struct MycDeliveryRecoveryEntropy([u8; 32]); 101 102 impl MycDeliveryRecoveryEntropy { 103 /// Wraps exact entropy from the caller's injected source. 104 #[must_use] 105 pub const fn from_injected_entropy(bytes: [u8; 32]) -> Self { 106 Self(bytes) 107 } 108 } 109 110 impl fmt::Debug for MycDeliveryRecoveryEntropy { 111 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 112 formatter.write_str("MycDeliveryRecoveryEntropy([redacted])") 113 } 114 } 115 116 #[derive(Clone, Copy, PartialEq, Eq)] 117 struct RecoveryCursor([u8; 32]); 118 119 /// Bounded, identity-free result of one restart-recovery transaction. 120 #[derive(Clone, Copy, PartialEq, Eq)] 121 pub struct MycDeliveryRecoveryReport { 122 examined_jobs: u32, 123 recovered_expired_attempts: u32, 124 finalized_jobs: u32, 125 promoted_discovery_generations: u32, 126 ready_targets: u32, 127 scheduled_targets: u32, 128 active_targets: u32, 129 continuation: Option<RecoveryCursor>, 130 } 131 132 impl MycDeliveryRecoveryReport { 133 #[must_use] 134 pub const fn examined_jobs(self) -> u32 { 135 self.examined_jobs 136 } 137 138 #[must_use] 139 pub const fn recovered_expired_attempts(self) -> u32 { 140 self.recovered_expired_attempts 141 } 142 143 #[must_use] 144 pub const fn finalized_jobs(self) -> u32 { 145 self.finalized_jobs 146 } 147 148 #[must_use] 149 pub const fn promoted_discovery_generations(self) -> u32 { 150 self.promoted_discovery_generations 151 } 152 153 #[must_use] 154 pub const fn ready_targets(self) -> u32 { 155 self.ready_targets 156 } 157 158 #[must_use] 159 pub const fn scheduled_targets(self) -> u32 { 160 self.scheduled_targets 161 } 162 163 #[must_use] 164 pub const fn active_targets(self) -> u32 { 165 self.active_targets 166 } 167 168 fn merge(&mut self, batch: Self) -> Result<(), RecoveryOperationError> { 169 self.examined_jobs = add(self.examined_jobs, batch.examined_jobs)?; 170 self.recovered_expired_attempts = add( 171 self.recovered_expired_attempts, 172 batch.recovered_expired_attempts, 173 )?; 174 self.finalized_jobs = add(self.finalized_jobs, batch.finalized_jobs)?; 175 self.promoted_discovery_generations = add( 176 self.promoted_discovery_generations, 177 batch.promoted_discovery_generations, 178 )?; 179 self.ready_targets = add(self.ready_targets, batch.ready_targets)?; 180 self.scheduled_targets = add(self.scheduled_targets, batch.scheduled_targets)?; 181 self.active_targets = add(self.active_targets, batch.active_targets)?; 182 self.continuation = batch.continuation; 183 Ok(()) 184 } 185 186 const fn empty() -> Self { 187 Self { 188 examined_jobs: 0, 189 recovered_expired_attempts: 0, 190 finalized_jobs: 0, 191 promoted_discovery_generations: 0, 192 ready_targets: 0, 193 scheduled_targets: 0, 194 active_targets: 0, 195 continuation: None, 196 } 197 } 198 } 199 200 impl fmt::Debug for MycDeliveryRecoveryReport { 201 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 202 formatter 203 .debug_struct("MycDeliveryRecoveryReport") 204 .field("examined_jobs", &self.examined_jobs) 205 .field( 206 "recovered_expired_attempts", 207 &self.recovered_expired_attempts, 208 ) 209 .field("finalized_jobs", &self.finalized_jobs) 210 .field( 211 "promoted_discovery_generations", 212 &self.promoted_discovery_generations, 213 ) 214 .field("ready_targets", &self.ready_targets) 215 .field("scheduled_targets", &self.scheduled_targets) 216 .field("active_targets", &self.active_targets) 217 .field("has_continuation", &self.continuation.is_some()) 218 .finish() 219 } 220 } 221 222 impl MycStateRepository<'_> { 223 /// Verifies the bounded global outbox relationships without mutation. 224 #[cfg(any(target_os = "linux", target_os = "macos"))] 225 pub(crate) async fn verify_delivery_invariants(&self) -> Result<(), MycStateRepositoryError> { 226 let expected = PersistedMetadata::from(self.expected()); 227 let outbox_maximum = self.expected().outbox_maximum(); 228 self.host() 229 .transaction(move |transaction| { 230 Box::pin(async move { 231 require_expected_metadata(transaction, &expected) 232 .await 233 .map_err(RecoveryOperationError::from)?; 234 verify_global_invariants(transaction, outbox_maximum).await 235 }) 236 }) 237 .await 238 .map_err(map_transaction_error) 239 } 240 241 /// Recovers startup delivery state in fixed bounded transactions without relay I/O. 242 /// 243 /// The later runtime owner must pause admission while invoking this method. 244 /// Pagination remains sealed inside the repository, time and entropy are 245 /// injected, and exact payload bytes and target membership never change. 246 pub async fn recover_delivery_state( 247 &self, 248 observed_at: MycDeliveryTimeUnixMs, 249 entropy: MycDeliveryRecoveryEntropy, 250 ) -> Result<MycDeliveryRecoveryReport, MycStateRepositoryError> { 251 let outbox_maximum = self.expected().outbox_maximum(); 252 let discovery_policy = self.expected().discovery_policies().cloned(); 253 let mut cursor = None; 254 let mut aggregate = MycDeliveryRecoveryReport::empty(); 255 loop { 256 let expected = PersistedMetadata::from(self.expected()); 257 let discovery_policy = discovery_policy.clone(); 258 let entropy_bytes = entropy.0; 259 let batch = self 260 .host() 261 .transaction(move |transaction| { 262 Box::pin(async move { 263 require_expected_metadata(transaction, &expected) 264 .await 265 .map_err(RecoveryOperationError::from)?; 266 recover_batch( 267 transaction, 268 observed_at, 269 &entropy_bytes, 270 cursor, 271 outbox_maximum, 272 discovery_policy.as_ref(), 273 ) 274 .await 275 }) 276 }) 277 .await 278 .map_err(map_transaction_error)?; 279 let next = batch.continuation; 280 aggregate 281 .merge(batch) 282 .map_err(|_| MycStateRepositoryError::new(MycStateRepositoryErrorKind::Binding))?; 283 let Some(next) = next else { 284 aggregate.continuation = None; 285 return Ok(aggregate); 286 }; 287 cursor = Some(next); 288 } 289 } 290 } 291 292 #[derive(Clone, Copy, Debug, PartialEq, Eq)] 293 enum RecoveryOperationError { 294 Binding, 295 Storage, 296 } 297 298 impl From<RepositoryOperationError> for RecoveryOperationError { 299 fn from(error: RepositoryOperationError) -> Self { 300 match error { 301 RepositoryOperationError::Binding => Self::Binding, 302 RepositoryOperationError::Storage => Self::Storage, 303 } 304 } 305 } 306 307 impl From<DeliveryOperationError> for RecoveryOperationError { 308 fn from(error: DeliveryOperationError) -> Self { 309 match error { 310 DeliveryOperationError::Binding => Self::Binding, 311 DeliveryOperationError::Storage => Self::Storage, 312 } 313 } 314 } 315 316 async fn recover_batch( 317 transaction: &mut ServiceSqliteTransaction<'_>, 318 observed_at: MycDeliveryTimeUnixMs, 319 entropy: &[u8; 32], 320 cursor: Option<RecoveryCursor>, 321 outbox_maximum: usize, 322 discovery_policy: Option<&MycDiscoveryPolicies>, 323 ) -> Result<MycDeliveryRecoveryReport, RecoveryOperationError> { 324 verify_global_invariants(transaction, outbox_maximum).await?; 325 let mut job_ids = read_job_ids(transaction, cursor).await?; 326 let has_more = job_ids.len() > MYC_DELIVERY_RECOVERY_BATCH_MAX_COUNT; 327 if has_more { 328 job_ids.truncate(MYC_DELIVERY_RECOVERY_BATCH_MAX_COUNT); 329 } 330 let continuation = has_more 331 .then(|| job_ids.last().copied().map(RecoveryCursor)) 332 .flatten(); 333 let mut report = MycDeliveryRecoveryReport { 334 examined_jobs: 0, 335 recovered_expired_attempts: 0, 336 finalized_jobs: 0, 337 promoted_discovery_generations: 0, 338 ready_targets: 0, 339 scheduled_targets: 0, 340 active_targets: 0, 341 continuation, 342 }; 343 for job_id in job_ids { 344 let before = read_job(transaction, MycDeliveryJobId::from_persisted(job_id)) 345 .await? 346 .ok_or(RecoveryOperationError::Binding)?; 347 match before.source_kind() { 348 MycDeliverySourceKind::SignerResponse => { 349 verify_response_for_delivery_job(transaction, before.id()).await?; 350 } 351 MycDeliverySourceKind::DiscoveryHandler => { 352 verify_document_for_delivery_job(transaction, &before, discovery_policy).await?; 353 } 354 } 355 validate_attempt_histories(transaction, &before).await?; 356 let mut current = before.clone(); 357 for target in before.targets() { 358 let Some(attempt_id) = target.active_attempt_id() else { 359 continue; 360 }; 361 let attempts = read_attempts(transaction, before.id(), target.index()).await?; 362 let attempt = attempts 363 .last() 364 .filter(|attempt| attempt.id() == attempt_id) 365 .ok_or(RecoveryOperationError::Binding)?; 366 if observed_at > attempt.lease_expires_at() { 367 let jitter = recovery_jitter(¤t, attempt.id(), attempt.number(), entropy)?; 368 current = recover_expired( 369 transaction, 370 before.id(), 371 target.relay_id(), 372 attempt.id(), 373 jitter, 374 observed_at, 375 ) 376 .await?; 377 report.recovered_expired_attempts = increment(report.recovered_expired_attempts)?; 378 } 379 } 380 current = finalize_job_if_terminal(transaction, before.id(), observed_at).await?; 381 if !is_terminal(before.status()) && is_terminal(current.status()) { 382 report.finalized_jobs = increment(report.finalized_jobs)?; 383 } 384 if current.status() == MycDeliveryJobStatus::Delivered 385 && current.source_kind() == MycDeliverySourceKind::DiscoveryHandler 386 && promote_current_if_desired(transaction, current.id(), observed_at).await? 387 { 388 report.promoted_discovery_generations = 389 increment(report.promoted_discovery_generations)?; 390 } 391 let verified = read_job(transaction, current.id()) 392 .await? 393 .ok_or(RecoveryOperationError::Binding)?; 394 validate_attempt_histories(transaction, &verified).await?; 395 classify_targets(&verified, observed_at, &mut report)?; 396 report.examined_jobs = increment(report.examined_jobs)?; 397 } 398 Ok(report) 399 } 400 401 async fn verify_global_invariants( 402 transaction: &mut ServiceSqliteTransaction<'_>, 403 outbox_maximum: usize, 404 ) -> Result<(), RecoveryOperationError> { 405 let rows = sqlx::query(READ_INVARIANTS_SQL) 406 .fetch_all(&mut *transaction) 407 .await 408 .map_err(|_| RecoveryOperationError::Storage)?; 409 if rows.len() != 1 { 410 return Err(RecoveryOperationError::Binding); 411 } 412 let row = &rows[0]; 413 for column in [ 414 "invalid_sources", 415 "orphan_responses", 416 "orphan_documents", 417 "orphan_targets", 418 "orphan_attempts", 419 ] { 420 if row 421 .try_get::<i64, _>(column) 422 .map_err(|_| RecoveryOperationError::Binding)? 423 != 0 424 { 425 return Err(RecoveryOperationError::Binding); 426 } 427 } 428 let active_jobs = row 429 .try_get::<i64, _>("active_jobs") 430 .map_err(|_| RecoveryOperationError::Binding)?; 431 usize::try_from(active_jobs) 432 .ok() 433 .filter(|count| *count <= outbox_maximum) 434 .map(|_| ()) 435 .ok_or(RecoveryOperationError::Binding) 436 } 437 438 async fn read_job_ids( 439 transaction: &mut ServiceSqliteTransaction<'_>, 440 cursor: Option<RecoveryCursor>, 441 ) -> Result<Vec<[u8; 32]>, RecoveryOperationError> { 442 let query = match cursor { 443 Some(cursor) => sqlx::query(READ_NEXT_JOB_IDS_SQL).bind(cursor.0.as_slice()), 444 None => sqlx::query(READ_FIRST_JOB_IDS_SQL), 445 }; 446 let rows = query 447 .fetch_all(&mut *transaction) 448 .await 449 .map_err(|_| RecoveryOperationError::Storage)?; 450 if rows.len() > MYC_DELIVERY_RECOVERY_BATCH_MAX_COUNT + 1 { 451 return Err(RecoveryOperationError::Binding); 452 } 453 rows.iter() 454 .map(|row| { 455 let bytes = row 456 .try_get::<Option<Vec<u8>>, _>("job_id") 457 .map_err(|_| RecoveryOperationError::Binding)? 458 .ok_or(RecoveryOperationError::Binding)?; 459 bytes 460 .try_into() 461 .map_err(|_| RecoveryOperationError::Binding) 462 }) 463 .collect() 464 } 465 466 async fn validate_attempt_histories( 467 transaction: &mut ServiceSqliteTransaction<'_>, 468 job: &MycDeliveryJobRecord, 469 ) -> Result<(), RecoveryOperationError> { 470 for target in job.targets() { 471 let attempts = read_attempts(transaction, job.id(), target.index()).await?; 472 if usize::try_from(target.attempt_count()) != Ok(attempts.len()) 473 || attempts 474 .iter() 475 .enumerate() 476 .any(|(index, attempt)| usize::try_from(attempt.number()) != Ok(index + 1)) 477 { 478 return Err(RecoveryOperationError::Binding); 479 } 480 let active = matches!( 481 target.status(), 482 MycDeliveryTargetStatus::Leased | MycDeliveryTargetStatus::Submitted 483 ); 484 if active { 485 let attempt = attempts.last().ok_or(RecoveryOperationError::Binding)?; 486 let expected_status = match target.status() { 487 MycDeliveryTargetStatus::Leased => MycDeliveryAttemptStatus::Leased, 488 MycDeliveryTargetStatus::Submitted => MycDeliveryAttemptStatus::Submitted, 489 _ => return Err(RecoveryOperationError::Binding), 490 }; 491 if target.active_attempt_id() != Some(attempt.id()) 492 || attempt.status() != expected_status 493 { 494 return Err(RecoveryOperationError::Binding); 495 } 496 } else if attempts.last().is_some_and(|attempt| { 497 matches!( 498 attempt.status(), 499 MycDeliveryAttemptStatus::Leased | MycDeliveryAttemptStatus::Submitted 500 ) 501 }) { 502 return Err(RecoveryOperationError::Binding); 503 } 504 } 505 Ok(()) 506 } 507 508 fn recovery_jitter( 509 job: &MycDeliveryJobRecord, 510 attempt_id: MycDeliveryAttemptId, 511 attempt_number: u32, 512 entropy: &[u8; 32], 513 ) -> Result<MycDeliveryRetryJitter, RecoveryOperationError> { 514 let exponent = attempt_number.saturating_sub(1).min(31); 515 let factor = 1_u64.checked_shl(exponent).unwrap_or(u64::MAX); 516 let maximum = job 517 .initial_backoff_ms() 518 .saturating_mul(factor) 519 .min(job.maximum_backoff_ms()); 520 let mut hasher = Sha256::new(); 521 hasher.update(RECOVERY_JITTER_DOMAIN); 522 hasher.update(entropy); 523 hasher.update(attempt_id.as_bytes()); 524 let digest: [u8; 32] = hasher.finalize().into(); 525 let raw = u64::from_be_bytes(digest[..8].try_into().expect("fixed SHA-256 prefix")); 526 let value = raw 527 % maximum 528 .checked_add(1) 529 .ok_or(RecoveryOperationError::Binding)?; 530 MycDeliveryRetryJitter::new(value).map_err(|_| RecoveryOperationError::Binding) 531 } 532 533 fn classify_targets( 534 job: &MycDeliveryJobRecord, 535 observed_at: MycDeliveryTimeUnixMs, 536 report: &mut MycDeliveryRecoveryReport, 537 ) -> Result<(), RecoveryOperationError> { 538 for target in job.targets() { 539 match target.status() { 540 MycDeliveryTargetStatus::Pending => { 541 report.ready_targets = increment(report.ready_targets)?; 542 } 543 MycDeliveryTargetStatus::Retryable | MycDeliveryTargetStatus::Unknown => { 544 if target 545 .next_attempt_at() 546 .is_some_and(|next| next > observed_at) 547 { 548 report.scheduled_targets = increment(report.scheduled_targets)?; 549 } else if target.next_attempt_at().is_some() { 550 report.ready_targets = increment(report.ready_targets)?; 551 } 552 } 553 MycDeliveryTargetStatus::Leased | MycDeliveryTargetStatus::Submitted => { 554 report.active_targets = increment(report.active_targets)?; 555 } 556 MycDeliveryTargetStatus::Delivered | MycDeliveryTargetStatus::Exhausted => {} 557 } 558 } 559 Ok(()) 560 } 561 562 const fn is_terminal(status: MycDeliveryJobStatus) -> bool { 563 matches!( 564 status, 565 MycDeliveryJobStatus::Delivered 566 | MycDeliveryJobStatus::Failed 567 | MycDeliveryJobStatus::Unknown 568 ) 569 } 570 571 fn increment(value: u32) -> Result<u32, RecoveryOperationError> { 572 value.checked_add(1).ok_or(RecoveryOperationError::Binding) 573 } 574 575 fn add(left: u32, right: u32) -> Result<u32, RecoveryOperationError> { 576 left.checked_add(right) 577 .ok_or(RecoveryOperationError::Binding) 578 } 579 580 fn map_transaction_error( 581 error: ServiceSqliteTransactionError<RecoveryOperationError>, 582 ) -> MycStateRepositoryError { 583 if error.kind() == ServiceSqliteTransactionErrorKind::CommitOutcomeUnknown { 584 return MycStateRepositoryError::new(MycStateRepositoryErrorKind::CommitOutcomeUnknown); 585 } 586 let kind = match error.operation_error() { 587 Some(RecoveryOperationError::Binding) => MycStateRepositoryErrorKind::Binding, 588 Some(RecoveryOperationError::Storage) | None => MycStateRepositoryErrorKind::Transaction, 589 }; 590 MycStateRepositoryError::new(kind) 591 }