publication_execution.rs (71086B)
1 //! Durable exact-byte publication claims, outcomes, retry, and recovery. 2 3 use core::fmt; 4 use std::error::Error; 5 6 use radroots_service_sqlite::{ 7 ServiceSqliteTransaction, ServiceSqliteTransactionError, ServiceSqliteTransactionErrorKind, 8 }; 9 use radroots_transport::BoxFuture; 10 use sha2::{Digest as _, Sha256}; 11 use sqlx::Row as _; 12 13 use crate::{ 14 RHI_PUBLICATION_ATTEMPT_NUMBER_MAXIMUM, RHI_PUBLICATION_TARGET_ORDINAL_MAXIMUM, 15 RhiCommittedPublication, RhiJitterBoundMilliseconds, RhiPublicationAttemptEvidence, 16 RhiPublicationAttemptId, RhiPublicationAttemptOutcome, RhiPublicationAuthority, 17 RhiPublicationMode, RhiPublicationOutboxId, RhiPublicationOutboxRepository, 18 RhiPublicationTargetState, RhiPublicationUnixMilliseconds, RhiRuntimeAdapterErrorKind, 19 RhiStateHostMode, RhiTimeEntropyAdapters, 20 publication_attempt::derive_attempt_id, 21 publication_submission::{ReadError, read_committed}, 22 }; 23 24 /// Exact version of the durable publication-execution contract. 25 pub const RHI_PUBLICATION_EXECUTION_CONTRACT_VERSION: u32 = 1; 26 27 const LEASE_OWNER_BYTES: usize = 16; 28 const MAX_UNIX_MILLISECONDS: u64 = i64::MAX as u64; 29 const TARGET_SET_DOMAIN: &[u8] = b"radroots.rhi.publication_target_set.v1\0"; 30 const READ_OUTBOX_SQL: &str = r#"SELECT 31 CASE WHEN typeof(outbox_id) = 'blob' AND length(outbox_id) = 32 32 THEN outbox_id ELSE NULL END AS outbox_id, 33 CASE WHEN typeof(event_sha256) = 'blob' AND length(event_sha256) = 32 34 THEN event_sha256 ELSE NULL END AS event_sha256, 35 CASE WHEN typeof(publication_authority_sha256) = 'blob' 36 AND length(publication_authority_sha256) = 32 37 THEN publication_authority_sha256 ELSE NULL END AS publication_authority_sha256, 38 CASE WHEN typeof(target_set_sha256) = 'blob' AND length(target_set_sha256) = 32 39 THEN target_set_sha256 ELSE NULL END AS target_set_sha256, 40 target_count, required_target_count, max_attempts, 41 initial_backoff_ms, maximum_backoff_ms, attempt_deadline_ms, 42 length(CAST(state AS BLOB)) AS state_bytes, substr(state, 1, 9) AS state, 43 revision, next_attempt_unix_ms, 44 CASE WHEN lease_owner IS NULL THEN NULL ELSE length(lease_owner) END AS lease_owner_bytes, 45 CASE WHEN lease_owner IS NULL THEN NULL ELSE substr(lease_owner, 1, 17) END AS lease_owner, 46 lease_expires_unix_ms, created_at_unix_ms, updated_at_unix_ms 47 FROM publication_outbox 48 WHERE outbox_id = ? 49 LIMIT 2"#; 50 51 const READ_CLAIMABLE_OUTBOX_SQL: &str = r#"SELECT 52 CASE WHEN typeof(outbox_id) = 'blob' AND length(outbox_id) = 32 53 THEN outbox_id ELSE NULL END AS outbox_id, 54 CASE WHEN typeof(event_sha256) = 'blob' AND length(event_sha256) = 32 55 THEN event_sha256 ELSE NULL END AS event_sha256, 56 CASE WHEN typeof(publication_authority_sha256) = 'blob' 57 AND length(publication_authority_sha256) = 32 58 THEN publication_authority_sha256 ELSE NULL END AS publication_authority_sha256, 59 CASE WHEN typeof(target_set_sha256) = 'blob' AND length(target_set_sha256) = 32 60 THEN target_set_sha256 ELSE NULL END AS target_set_sha256, 61 target_count, required_target_count, max_attempts, 62 initial_backoff_ms, maximum_backoff_ms, attempt_deadline_ms, 63 length(CAST(state AS BLOB)) AS state_bytes, substr(state, 1, 9) AS state, 64 revision, next_attempt_unix_ms, 65 NULL AS lease_owner_bytes, NULL AS lease_owner, 66 lease_expires_unix_ms, created_at_unix_ms, updated_at_unix_ms 67 FROM publication_outbox 68 WHERE state = 'pending' AND next_attempt_unix_ms <= ? 69 ORDER BY next_attempt_unix_ms, created_at_unix_ms, outbox_id 70 LIMIT 1"#; 71 72 const READ_EXPIRED_OUTBOX_SQL: &str = r#"SELECT 73 CASE WHEN typeof(outbox_id) = 'blob' AND length(outbox_id) = 32 74 THEN outbox_id ELSE NULL END AS outbox_id, 75 CASE WHEN typeof(event_sha256) = 'blob' AND length(event_sha256) = 32 76 THEN event_sha256 ELSE NULL END AS event_sha256, 77 CASE WHEN typeof(publication_authority_sha256) = 'blob' 78 AND length(publication_authority_sha256) = 32 79 THEN publication_authority_sha256 ELSE NULL END AS publication_authority_sha256, 80 CASE WHEN typeof(target_set_sha256) = 'blob' AND length(target_set_sha256) = 32 81 THEN target_set_sha256 ELSE NULL END AS target_set_sha256, 82 target_count, required_target_count, max_attempts, 83 initial_backoff_ms, maximum_backoff_ms, attempt_deadline_ms, 84 length(CAST(state AS BLOB)) AS state_bytes, substr(state, 1, 9) AS state, 85 revision, next_attempt_unix_ms, 86 length(lease_owner) AS lease_owner_bytes, substr(lease_owner, 1, 17) AS lease_owner, 87 lease_expires_unix_ms, created_at_unix_ms, updated_at_unix_ms 88 FROM publication_outbox 89 WHERE state = 'leased' AND lease_expires_unix_ms <= ? 90 ORDER BY lease_expires_unix_ms, created_at_unix_ms, outbox_id 91 LIMIT 1"#; 92 93 const READ_TARGETS_SQL: &str = r#"SELECT target_ordinal, 94 length(CAST(relay_id AS BLOB)) AS relay_id_bytes, substr(relay_id, 1, 65) AS relay_id, 95 required, 96 length(CAST(state AS BLOB)) AS state_bytes, substr(state, 1, 14) AS state, 97 revision, attempt_count, next_attempt_unix_ms, 98 CASE WHEN last_attempt_id IS NULL THEN NULL ELSE length(last_attempt_id) END 99 AS last_attempt_id_bytes, 100 CASE WHEN last_attempt_id IS NULL THEN NULL ELSE substr(last_attempt_id, 1, 33) END 101 AS last_attempt_id, 102 updated_at_unix_ms 103 FROM publication_targets 104 WHERE outbox_id = ? 105 ORDER BY target_ordinal 106 LIMIT 33"#; 107 108 const CLAIM_OUTBOX_SQL: &str = r#"UPDATE publication_outbox 109 SET state = 'leased', revision = revision + 1, 110 next_attempt_unix_ms = NULL, lease_owner = ?, lease_expires_unix_ms = ?, 111 updated_at_unix_ms = ? 112 WHERE outbox_id = ? AND revision = ? AND state = 'pending' 113 AND next_attempt_unix_ms <= ? AND updated_at_unix_ms <= ?"#; 114 115 const PREPARE_TARGET_SQL: &str = r#"UPDATE publication_targets 116 SET state = 'submitted', revision = revision + 1, 117 attempt_count = attempt_count + 1, next_attempt_unix_ms = NULL, 118 last_attempt_id = ?, updated_at_unix_ms = ? 119 WHERE outbox_id = ? AND target_ordinal = ? AND revision = ? 120 AND state IN ('pending', 'failed', 'rate_limited', 'unknown') 121 AND attempt_count < ? AND next_attempt_unix_ms <= ? 122 AND updated_at_unix_ms <= ?"#; 123 124 const INSERT_ATTEMPT_SQL: &str = r#"INSERT INTO publication_attempts ( 125 attempt_id, outbox_id, target_ordinal, attempt_number, event_sha256, 126 lease_owner, started_at_unix_ms, finished_at_unix_ms, outcome, result_code 127 ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)"#; 128 129 const READ_ATTEMPT_SQL: &str = r#"SELECT 130 CASE WHEN typeof(attempt_id) = 'blob' AND length(attempt_id) = 32 131 THEN attempt_id ELSE NULL END AS attempt_id, 132 CASE WHEN typeof(outbox_id) = 'blob' AND length(outbox_id) = 32 133 THEN outbox_id ELSE NULL END AS outbox_id, 134 target_ordinal, attempt_number, 135 CASE WHEN typeof(event_sha256) = 'blob' AND length(event_sha256) = 32 136 THEN event_sha256 ELSE NULL END AS event_sha256, 137 CASE WHEN typeof(lease_owner) = 'blob' AND length(lease_owner) = 16 138 THEN lease_owner ELSE NULL END AS lease_owner, 139 started_at_unix_ms, finished_at_unix_ms, 140 length(CAST(outcome AS BLOB)) AS outcome_bytes, substr(outcome, 1, 14) AS outcome, 141 length(CAST(result_code AS BLOB)) AS result_code_bytes, 142 substr(result_code, 1, 65) AS result_code 143 FROM publication_attempts 144 WHERE attempt_id = ? 145 LIMIT 2"#; 146 147 const UPDATE_TARGET_OUTCOME_SQL: &str = r#"UPDATE publication_targets 148 SET state = ?, revision = revision + 1, next_attempt_unix_ms = ?, 149 updated_at_unix_ms = ? 150 WHERE outbox_id = ? AND target_ordinal = ? AND revision = ? 151 AND state = 'submitted' AND attempt_count = ? AND last_attempt_id = ?"#; 152 153 const UPDATE_OUTBOX_AFTER_ATTEMPT_SQL: &str = r#"UPDATE publication_outbox 154 SET state = ?, revision = revision + 1, next_attempt_unix_ms = ?, 155 lease_owner = NULL, lease_expires_unix_ms = NULL, updated_at_unix_ms = ? 156 WHERE outbox_id = ? AND revision = ? AND state = 'leased' 157 AND lease_owner = ? AND lease_expires_unix_ms = ? 158 AND lease_expires_unix_ms > ? AND updated_at_unix_ms <= ?"#; 159 160 const UPDATE_OUTBOX_RECOVERY_SQL: &str = r#"UPDATE publication_outbox 161 SET state = ?, revision = revision + 1, next_attempt_unix_ms = ?, 162 lease_owner = NULL, lease_expires_unix_ms = NULL, updated_at_unix_ms = ? 163 WHERE outbox_id = ? AND revision = ? AND state = 'leased' 164 AND lease_owner = ? AND lease_expires_unix_ms = ? 165 AND lease_expires_unix_ms <= ? AND updated_at_unix_ms <= ?"#; 166 167 /// Stable durable outbox lifecycle state. 168 #[derive(Clone, Copy, Debug, PartialEq, Eq, Hash)] 169 pub enum RhiPublicationOutboxState { 170 Pending, 171 Leased, 172 Complete, 173 Blocked, 174 } 175 176 impl RhiPublicationOutboxState { 177 /// Returns the exact machine-contract spelling. 178 #[must_use] 179 pub const fn code(self) -> &'static str { 180 match self { 181 Self::Pending => "pending", 182 Self::Leased => "leased", 183 Self::Complete => "complete", 184 Self::Blocked => "blocked", 185 } 186 } 187 } 188 189 /// Stable source-free durable publication-execution failure class. 190 #[derive(Clone, Copy, Debug, PartialEq, Eq)] 191 pub enum RhiPublicationExecutionErrorKind { 192 InvalidMode, 193 InvalidInput, 194 NotReady, 195 LeaseLost, 196 Invariant, 197 ClockUnavailable, 198 EntropyUnavailable, 199 Storage, 200 CommitOutcomeUnknown, 201 } 202 203 impl RhiPublicationExecutionErrorKind { 204 /// Returns the stable machine-readable failure code. 205 #[must_use] 206 pub const fn code(self) -> &'static str { 207 match self { 208 Self::InvalidMode => "publication_execution_mode_invalid", 209 Self::InvalidInput => "publication_execution_input_invalid", 210 Self::NotReady => "publication_execution_not_ready", 211 Self::LeaseLost => "publication_execution_lease_lost", 212 Self::Invariant => "publication_execution_invariant_failed", 213 Self::ClockUnavailable => "publication_execution_clock_unavailable", 214 Self::EntropyUnavailable => "publication_execution_entropy_unavailable", 215 Self::Storage => "publication_execution_storage_failed", 216 Self::CommitOutcomeUnknown => "publication_execution_commit_outcome_unknown", 217 } 218 } 219 } 220 221 /// Redacted source-free durable publication-execution failure. 222 #[derive(Clone, Copy, PartialEq, Eq)] 223 pub struct RhiPublicationExecutionError { 224 kind: RhiPublicationExecutionErrorKind, 225 } 226 227 impl RhiPublicationExecutionError { 228 const fn new(kind: RhiPublicationExecutionErrorKind) -> Self { 229 Self { kind } 230 } 231 232 /// Returns the stable failure class. 233 #[must_use] 234 pub const fn kind(self) -> RhiPublicationExecutionErrorKind { 235 self.kind 236 } 237 238 /// Returns the stable machine-readable failure code. 239 #[must_use] 240 pub const fn code(self) -> &'static str { 241 self.kind.code() 242 } 243 } 244 245 impl fmt::Display for RhiPublicationExecutionError { 246 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 247 formatter.write_str(match self.kind { 248 RhiPublicationExecutionErrorKind::InvalidMode => { 249 "RHI publication execution requires writable state" 250 } 251 RhiPublicationExecutionErrorKind::InvalidInput => { 252 "RHI publication execution input is invalid" 253 } 254 RhiPublicationExecutionErrorKind::NotReady => "RHI publication work is not ready", 255 RhiPublicationExecutionErrorKind::LeaseLost => { 256 "RHI publication lease is no longer authoritative" 257 } 258 RhiPublicationExecutionErrorKind::Invariant => "RHI publication state invariant failed", 259 RhiPublicationExecutionErrorKind::ClockUnavailable => { 260 "RHI publication clock is unavailable" 261 } 262 RhiPublicationExecutionErrorKind::EntropyUnavailable => { 263 "RHI publication retry entropy is unavailable" 264 } 265 RhiPublicationExecutionErrorKind::Storage => "RHI publication transaction failed", 266 RhiPublicationExecutionErrorKind::CommitOutcomeUnknown => { 267 "RHI publication commit outcome is unknown" 268 } 269 }) 270 } 271 } 272 273 impl fmt::Debug for RhiPublicationExecutionError { 274 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 275 formatter 276 .debug_struct("RhiPublicationExecutionError") 277 .field("kind", &self.kind) 278 .finish() 279 } 280 } 281 282 impl Error for RhiPublicationExecutionError {} 283 284 /// Stable process-local owner token for one compare-and-swap publication lease. 285 #[derive(Clone, Copy, PartialEq, Eq, Hash)] 286 pub struct RhiPublicationLeaseOwner([u8; LEASE_OWNER_BYTES]); 287 288 impl RhiPublicationLeaseOwner { 289 /// Validates one injected nonzero lease-owner identity. 290 pub fn from_bytes( 291 bytes: [u8; LEASE_OWNER_BYTES], 292 ) -> Result<Self, RhiPublicationExecutionError> { 293 if bytes.iter().all(|byte| *byte == 0) { 294 return Err(RhiPublicationExecutionError::new( 295 RhiPublicationExecutionErrorKind::InvalidInput, 296 )); 297 } 298 Ok(Self(bytes)) 299 } 300 } 301 302 impl fmt::Debug for RhiPublicationLeaseOwner { 303 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 304 formatter.write_str("RhiPublicationLeaseOwner([redacted])") 305 } 306 } 307 308 /// Bounded caller-injected retry delay in whole milliseconds. 309 #[derive(Clone, Copy, Debug, PartialEq, Eq, PartialOrd, Ord)] 310 pub struct RhiPublicationRetryDelayMilliseconds(u64); 311 312 impl RhiPublicationRetryDelayMilliseconds { 313 /// Validates a delay against the absolute publication backoff ceiling. 314 pub fn new(value: u64) -> Result<Self, RhiPublicationExecutionError> { 315 if value > 3_600_000 { 316 return Err(RhiPublicationExecutionError::new( 317 RhiPublicationExecutionErrorKind::InvalidInput, 318 )); 319 } 320 Ok(Self(value)) 321 } 322 323 /// Returns the exact delay. 324 #[must_use] 325 pub const fn get(self) -> u64 { 326 self.0 327 } 328 } 329 330 /// Non-forgeable compare-and-swap authority for one claimed outbox. 331 #[derive(Clone, Copy, PartialEq, Eq)] 332 pub struct RhiPublicationLease { 333 outbox: OutboxRecord, 334 owner: RhiPublicationLeaseOwner, 335 expires_at: RhiPublicationUnixMilliseconds, 336 } 337 338 impl RhiPublicationLease { 339 /// Returns the exact claimed outbox identity. 340 #[must_use] 341 pub const fn outbox_id(self) -> RhiPublicationOutboxId { 342 self.outbox.id 343 } 344 345 /// Returns the exact lease expiry. 346 #[must_use] 347 pub const fn expires_at(self) -> RhiPublicationUnixMilliseconds { 348 self.expires_at 349 } 350 } 351 352 impl fmt::Debug for RhiPublicationLease { 353 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 354 formatter 355 .debug_struct("RhiPublicationLease") 356 .field("outbox", &"[redacted]") 357 .field("revision", &self.outbox.revision) 358 .field("expires_at", &self.expires_at) 359 .finish() 360 } 361 } 362 363 /// Sealed exact-byte remote submission prepared only after durable Submitted. 364 #[must_use = "prepared publication must be executed exactly or left for unknown recovery"] 365 pub struct RhiPreparedPublicationAttempt { 366 lease: RhiPublicationLease, 367 target: TargetRecord, 368 attempt_id: RhiPublicationAttemptId, 369 started_at: RhiPublicationUnixMilliseconds, 370 deadline_at: RhiPublicationUnixMilliseconds, 371 publication: RhiCommittedPublication, 372 } 373 374 impl RhiPreparedPublicationAttempt { 375 /// Returns the exact attempt identity. 376 #[must_use] 377 pub const fn attempt_id(&self) -> RhiPublicationAttemptId { 378 self.attempt_id 379 } 380 381 /// Returns the stable configured relay identity without an endpoint or secret. 382 #[must_use] 383 pub fn relay_id(&self) -> &str { 384 &self.target.relay_id 385 } 386 387 /// Returns the exact committed signed-event bytes with no transformation. 388 #[must_use] 389 pub const fn exact_signed_event_bytes(&self) -> &[u8] { 390 self.publication.exact_signed_event_bytes() 391 } 392 393 /// Returns the absolute attempt deadline. 394 #[must_use] 395 pub const fn deadline_at(&self) -> RhiPublicationUnixMilliseconds { 396 self.deadline_at 397 } 398 399 /// Returns the one-based durable attempt number. 400 #[must_use] 401 pub const fn attempt_number(&self) -> u16 { 402 self.target.attempt_count 403 } 404 405 fn retry_upper_bound(&self) -> u64 { 406 retry_upper_bound(self.lease.outbox, self.target.attempt_count) 407 } 408 } 409 410 impl fmt::Debug for RhiPreparedPublicationAttempt { 411 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 412 formatter 413 .debug_struct("RhiPreparedPublicationAttempt") 414 .field("identity", &"[redacted]") 415 .field("target_ordinal", &self.target.ordinal) 416 .field("attempt_number", &self.target.attempt_count) 417 .field("deadline_at", &self.deadline_at) 418 .finish() 419 } 420 } 421 422 /// Closed exact-byte transport boundary used by the durable executor. 423 /// 424 /// The adapter receives the original committed payload and must submit that 425 /// byte slice unchanged. Dropping the future after durable preparation leaves 426 /// Submitted evidence; expired-lease recovery records Unknown before retry. 427 pub trait RhiExactPublicationSink: Send + Sync { 428 /// Submits one exact prepared payload and returns only a closed observation. 429 fn submit_exact<'a>( 430 &'a self, 431 attempt: &'a RhiPreparedPublicationAttempt, 432 ) -> BoxFuture<'a, RhiPublicationAttemptOutcome>; 433 } 434 435 /// Confirmed durable result for one exact publication attempt. 436 #[derive(Clone, Copy, Debug, PartialEq, Eq)] 437 pub struct RhiPublicationAttemptCommit { 438 outbox_id: RhiPublicationOutboxId, 439 attempt_id: RhiPublicationAttemptId, 440 target_ordinal: u8, 441 attempt_number: u16, 442 outcome: RhiPublicationAttemptOutcome, 443 target_state: RhiPublicationTargetState, 444 outbox_state: RhiPublicationOutboxState, 445 } 446 447 impl RhiPublicationAttemptCommit { 448 #[must_use] 449 pub const fn outbox_id(self) -> RhiPublicationOutboxId { 450 self.outbox_id 451 } 452 453 #[must_use] 454 pub const fn attempt_id(self) -> RhiPublicationAttemptId { 455 self.attempt_id 456 } 457 458 #[must_use] 459 pub const fn target_ordinal(self) -> u8 { 460 self.target_ordinal 461 } 462 463 #[must_use] 464 pub const fn attempt_number(self) -> u16 { 465 self.attempt_number 466 } 467 468 #[must_use] 469 pub const fn outcome(self) -> RhiPublicationAttemptOutcome { 470 self.outcome 471 } 472 473 #[must_use] 474 pub const fn target_state(self) -> RhiPublicationTargetState { 475 self.target_state 476 } 477 478 #[must_use] 479 pub const fn outbox_state(self) -> RhiPublicationOutboxState { 480 self.outbox_state 481 } 482 } 483 484 impl RhiPublicationOutboxRepository<'_> { 485 /// Claims the oldest due outbox using one injected owner and wall time. 486 pub async fn claim_next_publication( 487 &self, 488 owner: RhiPublicationLeaseOwner, 489 now: RhiPublicationUnixMilliseconds, 490 authority: &RhiPublicationAuthority, 491 ) -> Result<Option<RhiPublicationLease>, RhiPublicationExecutionError> { 492 require_writable(self)?; 493 let Some(authority) = AuthorityBinding::from_authority(authority) else { 494 return Ok(None); 495 }; 496 self.host() 497 .sqlite_host() 498 .transaction(move |transaction| { 499 Box::pin(async move { claim_next(transaction, owner, now, authority).await }) 500 }) 501 .await 502 .map_err(map_transaction_error) 503 } 504 505 /// Persists Submitted before exposing exact bytes to remote I/O. 506 pub async fn prepare_next_publication_target( 507 &self, 508 lease: RhiPublicationLease, 509 started_at: RhiPublicationUnixMilliseconds, 510 ) -> Result<RhiPreparedPublicationAttempt, RhiPublicationExecutionError> { 511 require_writable(self)?; 512 self.host() 513 .sqlite_host() 514 .transaction(move |transaction| { 515 Box::pin(async move { prepare_next(transaction, lease, started_at).await }) 516 }) 517 .await 518 .map_err(map_transaction_error) 519 } 520 521 /// Commits one closed observed outcome by exact compare-and-swap. 522 pub async fn record_publication_outcome( 523 &self, 524 prepared: &RhiPreparedPublicationAttempt, 525 finished_at: RhiPublicationUnixMilliseconds, 526 outcome: RhiPublicationAttemptOutcome, 527 retry_delay: RhiPublicationRetryDelayMilliseconds, 528 ) -> Result<RhiPublicationAttemptCommit, RhiPublicationExecutionError> { 529 require_writable(self)?; 530 let input = RecordInput::from_prepared(prepared, finished_at, outcome, retry_delay)?; 531 self.host() 532 .sqlite_host() 533 .transaction(move |transaction| { 534 Box::pin(async move { record_outcome(transaction, input).await }) 535 }) 536 .await 537 .map_err(map_transaction_error) 538 } 539 540 /// Recovers at most one expired lease and persists Unknown for Submitted work. 541 pub async fn recover_one_expired_publication( 542 &self, 543 adapters: &RhiTimeEntropyAdapters, 544 now: RhiPublicationUnixMilliseconds, 545 authority: &RhiPublicationAuthority, 546 ) -> Result<bool, RhiPublicationExecutionError> { 547 require_writable(self)?; 548 let Some(authority) = AuthorityBinding::from_authority(authority) else { 549 return Ok(false); 550 }; 551 let candidate = self 552 .host() 553 .sqlite_host() 554 .transaction(move |transaction| { 555 Box::pin(async move { read_expired_candidate(transaction, now, authority).await }) 556 }) 557 .await 558 .map_err(map_transaction_error)?; 559 let Some(candidate) = candidate else { 560 return Ok(false); 561 }; 562 let cap = candidate.retry_upper_bound; 563 let delay = sample_retry_delay(adapters, cap)?; 564 self.host() 565 .sqlite_host() 566 .transaction(move |transaction| { 567 Box::pin(async move { recover_expired(transaction, candidate, now, delay).await }) 568 }) 569 .await 570 .map_err(map_transaction_error) 571 } 572 573 /// Executes at most one exact-byte attempt through the injected sink. 574 /// 575 /// Cancellation before a claim has no effect. Cancellation after durable 576 /// preparation leaves Submitted state; lease-expiry recovery records 577 /// Unknown and schedules the exact same committed bytes without rebuilding. 578 pub async fn execute_next_publication( 579 &self, 580 owner: RhiPublicationLeaseOwner, 581 adapters: &RhiTimeEntropyAdapters, 582 sink: &dyn RhiExactPublicationSink, 583 authority: &RhiPublicationAuthority, 584 ) -> Result<Option<RhiPublicationAttemptCommit>, RhiPublicationExecutionError> { 585 let now = publication_now(adapters)?; 586 self.recover_one_expired_publication(adapters, now, authority) 587 .await?; 588 let Some(lease) = self.claim_next_publication(owner, now, authority).await? else { 589 return Ok(None); 590 }; 591 let prepared = self.prepare_next_publication_target(lease, now).await?; 592 let outcome = match sink.submit_exact(&prepared).await { 593 RhiPublicationAttemptOutcome::Submitted => RhiPublicationAttemptOutcome::Unknown, 594 outcome => outcome, 595 }; 596 let finished_at = publication_now(adapters)?; 597 let retry_delay = if retryable_outcome(outcome) { 598 sample_retry_delay(adapters, prepared.retry_upper_bound())? 599 } else { 600 RhiPublicationRetryDelayMilliseconds(0) 601 }; 602 self.record_publication_outcome(&prepared, finished_at, outcome, retry_delay) 603 .await 604 .map(Some) 605 } 606 } 607 608 fn require_writable( 609 repository: &RhiPublicationOutboxRepository<'_>, 610 ) -> Result<(), RhiPublicationExecutionError> { 611 if repository.host().mode() == RhiStateHostMode::ReadWriteExisting { 612 Ok(()) 613 } else { 614 Err(RhiPublicationExecutionError::new( 615 RhiPublicationExecutionErrorKind::InvalidMode, 616 )) 617 } 618 } 619 620 fn publication_now( 621 adapters: &RhiTimeEntropyAdapters, 622 ) -> Result<RhiPublicationUnixMilliseconds, RhiPublicationExecutionError> { 623 let value = adapters.now_utc_milliseconds().map_err(|error| { 624 RhiPublicationExecutionError::new(match error.kind() { 625 RhiRuntimeAdapterErrorKind::WallClockUnavailable => { 626 RhiPublicationExecutionErrorKind::ClockUnavailable 627 } 628 _ => RhiPublicationExecutionErrorKind::ClockUnavailable, 629 }) 630 })?; 631 RhiPublicationUnixMilliseconds::new(value).map_err(|_| { 632 RhiPublicationExecutionError::new(RhiPublicationExecutionErrorKind::ClockUnavailable) 633 }) 634 } 635 636 fn sample_retry_delay( 637 adapters: &RhiTimeEntropyAdapters, 638 cap: u64, 639 ) -> Result<RhiPublicationRetryDelayMilliseconds, RhiPublicationExecutionError> { 640 if cap == 0 { 641 return Ok(RhiPublicationRetryDelayMilliseconds(0)); 642 } 643 let bound = RhiJitterBoundMilliseconds::new(cap).map_err(|_| { 644 RhiPublicationExecutionError::new(RhiPublicationExecutionErrorKind::InvalidInput) 645 })?; 646 adapters 647 .sample_full_jitter(bound) 648 .map(|delay| RhiPublicationRetryDelayMilliseconds(delay.get())) 649 .map_err(|_| { 650 RhiPublicationExecutionError::new(RhiPublicationExecutionErrorKind::EntropyUnavailable) 651 }) 652 } 653 654 async fn claim_next( 655 transaction: &mut ServiceSqliteTransaction<'_>, 656 owner: RhiPublicationLeaseOwner, 657 now: RhiPublicationUnixMilliseconds, 658 authority: AuthorityBinding, 659 ) -> Result<Option<RhiPublicationLease>, OperationError> { 660 let rows = sqlx::query(READ_CLAIMABLE_OUTBOX_SQL) 661 .bind(i64_value(now.get())?) 662 .fetch_all(&mut *transaction) 663 .await 664 .map_err(|_| OperationError::Storage)?; 665 let Some(row) = exactly_zero_or_one(rows)? else { 666 return Ok(None); 667 }; 668 let outbox = decode_outbox(row)?; 669 validate_authority(outbox, authority)?; 670 validate_target_inventory(transaction, outbox).await?; 671 let expires = now 672 .get() 673 .checked_add(outbox.attempt_deadline_ms) 674 .filter(|value| *value <= MAX_UNIX_MILLISECONDS) 675 .ok_or(OperationError::InvalidInput)?; 676 let result = sqlx::query(CLAIM_OUTBOX_SQL) 677 .bind(owner.0.as_slice()) 678 .bind(i64_value(expires)?) 679 .bind(i64_value(now.get())?) 680 .bind(outbox.id.as_bytes().as_slice()) 681 .bind(i64_value(outbox.revision)?) 682 .bind(i64_value(now.get())?) 683 .bind(i64_value(now.get())?) 684 .execute(&mut *transaction) 685 .await 686 .map_err(|_| OperationError::Storage)?; 687 if result.rows_affected() != 1 { 688 return Err(OperationError::LeaseLost); 689 } 690 let leased = read_outbox(transaction, outbox.id) 691 .await? 692 .filter(|record| { 693 record.state == RhiPublicationOutboxState::Leased 694 && record.lease_owner == Some(owner) 695 && record.lease_expires == Some(RhiPublicationUnixMilliseconds(expires)) 696 }) 697 .ok_or(OperationError::Invariant)?; 698 Ok(Some(RhiPublicationLease { 699 outbox: leased, 700 owner, 701 expires_at: RhiPublicationUnixMilliseconds(expires), 702 })) 703 } 704 705 async fn prepare_next( 706 transaction: &mut ServiceSqliteTransaction<'_>, 707 lease: RhiPublicationLease, 708 started_at: RhiPublicationUnixMilliseconds, 709 ) -> Result<RhiPreparedPublicationAttempt, OperationError> { 710 validate_lease(transaction, lease, started_at, true).await?; 711 let targets = read_targets(transaction, lease.outbox.id).await?; 712 validate_targets(lease.outbox, &targets)?; 713 let candidate = targets 714 .into_iter() 715 .find(|target| target_is_due(target, lease.outbox.max_attempts, started_at)) 716 .ok_or(OperationError::NotReady)?; 717 let attempt_number = candidate 718 .attempt_count 719 .checked_add(1) 720 .filter(|value| *value <= lease.outbox.max_attempts) 721 .ok_or(OperationError::Invariant)?; 722 let attempt_id = derive_attempt_id( 723 lease.outbox.id, 724 lease.outbox.event_sha256, 725 candidate.ordinal, 726 attempt_number, 727 ); 728 let result = sqlx::query(PREPARE_TARGET_SQL) 729 .bind(attempt_id.as_bytes().as_slice()) 730 .bind(i64_value(started_at.get())?) 731 .bind(lease.outbox.id.as_bytes().as_slice()) 732 .bind(i64::from(candidate.ordinal)) 733 .bind(i64_value(candidate.revision)?) 734 .bind(i64::from(lease.outbox.max_attempts)) 735 .bind(i64_value(started_at.get())?) 736 .bind(i64_value(started_at.get())?) 737 .execute(&mut *transaction) 738 .await 739 .map_err(|_| OperationError::Storage)?; 740 if result.rows_affected() != 1 { 741 return Err(OperationError::LeaseLost); 742 } 743 let target = read_targets(transaction, lease.outbox.id) 744 .await? 745 .into_iter() 746 .find(|target| target.ordinal == candidate.ordinal) 747 .filter(|target| { 748 target.state == RhiPublicationTargetState::Submitted 749 && target.attempt_count == attempt_number 750 && target.last_attempt_id == Some(attempt_id) 751 }) 752 .ok_or(OperationError::Invariant)?; 753 let publication = read_committed(transaction, lease.outbox.id) 754 .await 755 .map_err(|error| match error { 756 ReadError::Storage => OperationError::Storage, 757 ReadError::NotFound | ReadError::Binding => OperationError::Invariant, 758 })?; 759 if publication.event_sha256() != &lease.outbox.event_sha256 { 760 return Err(OperationError::Invariant); 761 } 762 Ok(RhiPreparedPublicationAttempt { 763 lease, 764 target, 765 attempt_id, 766 started_at, 767 deadline_at: lease.expires_at, 768 publication, 769 }) 770 } 771 772 async fn record_outcome( 773 transaction: &mut ServiceSqliteTransaction<'_>, 774 input: RecordInput, 775 ) -> Result<RhiPublicationAttemptCommit, OperationError> { 776 if let Some(existing) = read_attempt(transaction, input.attempt_id).await? { 777 return reconcile_recorded(transaction, input, existing).await; 778 } 779 validate_lease(transaction, input.lease, input.finished_at, true).await?; 780 let target = read_targets(transaction, input.lease.outbox.id) 781 .await? 782 .into_iter() 783 .find(|target| target.ordinal == input.target_ordinal) 784 .ok_or(OperationError::Invariant)?; 785 if target.state != RhiPublicationTargetState::Submitted 786 || target.revision != input.target_revision 787 || target.attempt_count != input.attempt_number 788 || target.last_attempt_id != Some(input.attempt_id) 789 || target.updated_at != input.started_at 790 { 791 return Err(OperationError::LeaseLost); 792 } 793 insert_attempt(transaction, input).await?; 794 let (next_state, next_attempt) = target_outcome_schedule(input)?; 795 let result = sqlx::query(UPDATE_TARGET_OUTCOME_SQL) 796 .bind(next_state.code()) 797 .bind(optional_i64(next_attempt)?) 798 .bind(i64_value(input.finished_at.get())?) 799 .bind(input.lease.outbox.id.as_bytes().as_slice()) 800 .bind(i64::from(input.target_ordinal)) 801 .bind(i64_value(input.target_revision)?) 802 .bind(i64::from(input.attempt_number)) 803 .bind(input.attempt_id.as_bytes().as_slice()) 804 .execute(&mut *transaction) 805 .await 806 .map_err(|_| OperationError::Storage)?; 807 if result.rows_affected() != 1 { 808 return Err(OperationError::LeaseLost); 809 } 810 let targets = read_targets(transaction, input.lease.outbox.id).await?; 811 let disposition = disposition(input.lease.outbox, &targets)?; 812 update_outbox_after_attempt(transaction, input, disposition).await?; 813 Ok(RhiPublicationAttemptCommit { 814 outbox_id: input.lease.outbox.id, 815 attempt_id: input.attempt_id, 816 target_ordinal: input.target_ordinal, 817 attempt_number: input.attempt_number, 818 outcome: input.outcome, 819 target_state: next_state, 820 outbox_state: disposition.state, 821 }) 822 } 823 824 async fn read_expired_candidate( 825 transaction: &mut ServiceSqliteTransaction<'_>, 826 now: RhiPublicationUnixMilliseconds, 827 authority: AuthorityBinding, 828 ) -> Result<Option<RecoveryCandidate>, OperationError> { 829 let rows = sqlx::query(READ_EXPIRED_OUTBOX_SQL) 830 .bind(i64_value(now.get())?) 831 .fetch_all(&mut *transaction) 832 .await 833 .map_err(|_| OperationError::Storage)?; 834 let Some(row) = exactly_zero_or_one(rows)? else { 835 return Ok(None); 836 }; 837 let outbox = decode_outbox(row)?; 838 validate_authority(outbox, authority)?; 839 let targets = read_targets(transaction, outbox.id).await?; 840 validate_targets(outbox, &targets)?; 841 let submitted_attempts: Vec<_> = targets 842 .iter() 843 .filter(|target| target.state == RhiPublicationTargetState::Submitted) 844 .map(|target| target.attempt_count) 845 .collect(); 846 if submitted_attempts.len() > 1 { 847 return Err(OperationError::Invariant); 848 } 849 let retry_upper_bound = submitted_attempts 850 .first() 851 .copied() 852 .filter(|attempt| *attempt < outbox.max_attempts) 853 .map_or(0, |attempt| retry_upper_bound(outbox, attempt)); 854 Ok(Some(RecoveryCandidate { 855 outbox, 856 retry_upper_bound, 857 })) 858 } 859 860 async fn recover_expired( 861 transaction: &mut ServiceSqliteTransaction<'_>, 862 candidate: RecoveryCandidate, 863 now: RhiPublicationUnixMilliseconds, 864 delay: RhiPublicationRetryDelayMilliseconds, 865 ) -> Result<bool, OperationError> { 866 let current = read_outbox(transaction, candidate.outbox.id) 867 .await? 868 .filter(|current| *current == candidate.outbox) 869 .ok_or(OperationError::LeaseLost)?; 870 if current.state != RhiPublicationOutboxState::Leased 871 || current.lease_expires.is_none_or(|expires| expires > now) 872 { 873 return Err(OperationError::LeaseLost); 874 } 875 let owner = current.lease_owner.ok_or(OperationError::Invariant)?; 876 let mut targets = read_targets(transaction, current.id).await?; 877 validate_targets(current, &targets)?; 878 for target in targets 879 .iter_mut() 880 .filter(|target| target.state == RhiPublicationTargetState::Submitted) 881 { 882 let attempt_id = target.last_attempt_id.ok_or(OperationError::Invariant)?; 883 let publication = 884 read_committed(transaction, current.id) 885 .await 886 .map_err(|error| match error { 887 ReadError::Storage => OperationError::Storage, 888 ReadError::NotFound | ReadError::Binding => OperationError::Invariant, 889 })?; 890 let evidence = RhiPublicationAttemptEvidence::new( 891 &publication, 892 u32::from(target.ordinal), 893 target.attempt_count, 894 target.updated_at, 895 now, 896 RhiPublicationAttemptOutcome::Unknown, 897 ) 898 .map_err(|_| OperationError::Invariant)?; 899 if attempt_id != evidence.id() || read_attempt(transaction, attempt_id).await?.is_some() { 900 return Err(OperationError::Invariant); 901 } 902 let input = RecordInput { 903 lease: RhiPublicationLease { 904 outbox: current, 905 owner, 906 expires_at: current.lease_expires.ok_or(OperationError::Invariant)?, 907 }, 908 target_ordinal: target.ordinal, 909 target_revision: target.revision, 910 attempt_number: target.attempt_count, 911 attempt_id, 912 started_at: target.updated_at, 913 finished_at: now, 914 outcome: RhiPublicationAttemptOutcome::Unknown, 915 retry_delay: delay, 916 }; 917 insert_attempt(transaction, input).await?; 918 let (state, next) = target_outcome_schedule(input)?; 919 let result = sqlx::query(UPDATE_TARGET_OUTCOME_SQL) 920 .bind(state.code()) 921 .bind(optional_i64(next)?) 922 .bind(i64_value(now.get())?) 923 .bind(current.id.as_bytes().as_slice()) 924 .bind(i64::from(target.ordinal)) 925 .bind(i64_value(target.revision)?) 926 .bind(i64::from(target.attempt_count)) 927 .bind(attempt_id.as_bytes().as_slice()) 928 .execute(&mut *transaction) 929 .await 930 .map_err(|_| OperationError::Storage)?; 931 if result.rows_affected() != 1 { 932 return Err(OperationError::LeaseLost); 933 } 934 } 935 targets = read_targets(transaction, current.id).await?; 936 let disposition = disposition(current, &targets)?; 937 let result = sqlx::query(UPDATE_OUTBOX_RECOVERY_SQL) 938 .bind(disposition.state.code()) 939 .bind(optional_i64(disposition.next_attempt)?) 940 .bind(i64_value(now.get())?) 941 .bind(current.id.as_bytes().as_slice()) 942 .bind(i64_value(current.revision)?) 943 .bind(owner.0.as_slice()) 944 .bind(i64_value( 945 current 946 .lease_expires 947 .ok_or(OperationError::Invariant)? 948 .get(), 949 )?) 950 .bind(i64_value(now.get())?) 951 .bind(i64_value(now.get())?) 952 .execute(&mut *transaction) 953 .await 954 .map_err(|_| OperationError::Storage)?; 955 if result.rows_affected() != 1 { 956 return Err(OperationError::LeaseLost); 957 } 958 Ok(true) 959 } 960 961 async fn validate_lease( 962 transaction: &mut ServiceSqliteTransaction<'_>, 963 lease: RhiPublicationLease, 964 now: RhiPublicationUnixMilliseconds, 965 require_unexpired: bool, 966 ) -> Result<(), OperationError> { 967 let current = read_outbox(transaction, lease.outbox.id) 968 .await? 969 .filter(|current| *current == lease.outbox) 970 .ok_or(OperationError::LeaseLost)?; 971 if current.state != RhiPublicationOutboxState::Leased 972 || current.lease_owner != Some(lease.owner) 973 || current.lease_expires != Some(lease.expires_at) 974 || (require_unexpired && lease.expires_at <= now) 975 { 976 return Err(OperationError::LeaseLost); 977 } 978 Ok(()) 979 } 980 981 async fn insert_attempt( 982 transaction: &mut ServiceSqliteTransaction<'_>, 983 input: RecordInput, 984 ) -> Result<(), OperationError> { 985 let result = sqlx::query(INSERT_ATTEMPT_SQL) 986 .bind(input.attempt_id.as_bytes().as_slice()) 987 .bind(input.lease.outbox.id.as_bytes().as_slice()) 988 .bind(i64::from(input.target_ordinal)) 989 .bind(i64::from(input.attempt_number)) 990 .bind(input.lease.outbox.event_sha256.as_slice()) 991 .bind(input.lease.owner.0.as_slice()) 992 .bind(i64_value(input.started_at.get())?) 993 .bind(i64_value(input.finished_at.get())?) 994 .bind(input.outcome.code()) 995 .bind(input.outcome.code()) 996 .execute(&mut *transaction) 997 .await 998 .map_err(|_| OperationError::Storage)?; 999 if result.rows_affected() != 1 { 1000 return Err(OperationError::Storage); 1001 } 1002 Ok(()) 1003 } 1004 1005 async fn reconcile_recorded( 1006 transaction: &mut ServiceSqliteTransaction<'_>, 1007 input: RecordInput, 1008 existing: AttemptRecord, 1009 ) -> Result<RhiPublicationAttemptCommit, OperationError> { 1010 if existing != AttemptRecord::from_input(input) { 1011 return Err(OperationError::Invariant); 1012 } 1013 let outbox = read_outbox(transaction, input.lease.outbox.id) 1014 .await? 1015 .ok_or(OperationError::Invariant)?; 1016 if !same_outbox_identity(outbox, input.lease.outbox) 1017 || outbox.revision 1018 != input 1019 .lease 1020 .outbox 1021 .revision 1022 .checked_add(1) 1023 .ok_or(OperationError::Invariant)? 1024 || outbox.updated_at != input.finished_at 1025 || outbox.lease_owner.is_some() 1026 || outbox.lease_expires.is_some() 1027 { 1028 return Err(OperationError::Invariant); 1029 } 1030 let targets = read_targets(transaction, outbox.id).await?; 1031 validate_targets(outbox, &targets)?; 1032 let (expected_target_state, expected_target_schedule) = target_outcome_schedule(input)?; 1033 let expected_target_revision = input 1034 .target_revision 1035 .checked_add(1) 1036 .ok_or(OperationError::Invariant)?; 1037 let target = targets 1038 .iter() 1039 .find(|target| target.ordinal == input.target_ordinal) 1040 .filter(|target| { 1041 target.revision == expected_target_revision 1042 && target.attempt_count == input.attempt_number 1043 && target.last_attempt_id == Some(input.attempt_id) 1044 && target.state == expected_target_state 1045 && target.next_attempt == expected_target_schedule 1046 && target.updated_at == input.finished_at 1047 }) 1048 .ok_or(OperationError::Invariant)?; 1049 let expected_disposition = disposition(outbox, &targets)?; 1050 if outbox.state != expected_disposition.state 1051 || outbox.next_attempt != expected_disposition.next_attempt 1052 { 1053 return Err(OperationError::Invariant); 1054 } 1055 Ok(RhiPublicationAttemptCommit { 1056 outbox_id: outbox.id, 1057 attempt_id: input.attempt_id, 1058 target_ordinal: input.target_ordinal, 1059 attempt_number: input.attempt_number, 1060 outcome: input.outcome, 1061 target_state: target.state, 1062 outbox_state: outbox.state, 1063 }) 1064 } 1065 1066 fn same_outbox_identity(current: OutboxRecord, prior: OutboxRecord) -> bool { 1067 current.id == prior.id 1068 && current.event_sha256 == prior.event_sha256 1069 && current.authority_sha256 == prior.authority_sha256 1070 && current.target_set_sha256 == prior.target_set_sha256 1071 && current.target_count == prior.target_count 1072 && current.required_target_count == prior.required_target_count 1073 && current.max_attempts == prior.max_attempts 1074 && current.initial_backoff_ms == prior.initial_backoff_ms 1075 && current.maximum_backoff_ms == prior.maximum_backoff_ms 1076 && current.attempt_deadline_ms == prior.attempt_deadline_ms 1077 && current.created_at == prior.created_at 1078 } 1079 1080 async fn update_outbox_after_attempt( 1081 transaction: &mut ServiceSqliteTransaction<'_>, 1082 input: RecordInput, 1083 disposition: OutboxDisposition, 1084 ) -> Result<(), OperationError> { 1085 let result = sqlx::query(UPDATE_OUTBOX_AFTER_ATTEMPT_SQL) 1086 .bind(disposition.state.code()) 1087 .bind(optional_i64(disposition.next_attempt)?) 1088 .bind(i64_value(input.finished_at.get())?) 1089 .bind(input.lease.outbox.id.as_bytes().as_slice()) 1090 .bind(i64_value(input.lease.outbox.revision)?) 1091 .bind(input.lease.owner.0.as_slice()) 1092 .bind(i64_value(input.lease.expires_at.get())?) 1093 .bind(i64_value(input.finished_at.get())?) 1094 .bind(i64_value(input.finished_at.get())?) 1095 .execute(&mut *transaction) 1096 .await 1097 .map_err(|_| OperationError::Storage)?; 1098 if result.rows_affected() != 1 { 1099 return Err(OperationError::LeaseLost); 1100 } 1101 Ok(()) 1102 } 1103 1104 async fn read_outbox( 1105 transaction: &mut ServiceSqliteTransaction<'_>, 1106 id: RhiPublicationOutboxId, 1107 ) -> Result<Option<OutboxRecord>, OperationError> { 1108 let rows = sqlx::query(READ_OUTBOX_SQL) 1109 .bind(id.as_bytes().as_slice()) 1110 .fetch_all(&mut *transaction) 1111 .await 1112 .map_err(|_| OperationError::Storage)?; 1113 exactly_zero_or_one(rows)?.map(decode_outbox).transpose() 1114 } 1115 1116 async fn read_targets( 1117 transaction: &mut ServiceSqliteTransaction<'_>, 1118 id: RhiPublicationOutboxId, 1119 ) -> Result<Vec<TargetRecord>, OperationError> { 1120 sqlx::query(READ_TARGETS_SQL) 1121 .bind(id.as_bytes().as_slice()) 1122 .fetch_all(&mut *transaction) 1123 .await 1124 .map_err(|_| OperationError::Storage)? 1125 .into_iter() 1126 .map(decode_target) 1127 .collect() 1128 } 1129 1130 async fn validate_target_inventory( 1131 transaction: &mut ServiceSqliteTransaction<'_>, 1132 outbox: OutboxRecord, 1133 ) -> Result<(), OperationError> { 1134 let targets = read_targets(transaction, outbox.id).await?; 1135 validate_targets(outbox, &targets) 1136 } 1137 1138 fn validate_targets(outbox: OutboxRecord, targets: &[TargetRecord]) -> Result<(), OperationError> { 1139 if targets.len() != usize::from(outbox.target_count) 1140 || targets.len() 1141 > usize::try_from(RHI_PUBLICATION_TARGET_ORDINAL_MAXIMUM + 1) 1142 .map_err(|_| OperationError::Invariant)? 1143 || targets 1144 .iter() 1145 .enumerate() 1146 .any(|(index, target)| usize::from(target.ordinal) != index) 1147 || targets.iter().filter(|target| target.required).count() 1148 != usize::from(outbox.required_target_count) 1149 { 1150 return Err(OperationError::Invariant); 1151 } 1152 for (index, target) in targets.iter().enumerate() { 1153 if targets[index + 1..] 1154 .iter() 1155 .any(|other| other.relay_id == target.relay_id) 1156 { 1157 return Err(OperationError::Invariant); 1158 } 1159 } 1160 let mut digest = Sha256::new(); 1161 digest.update(TARGET_SET_DOMAIN); 1162 digest.update(u32::from(outbox.target_count).to_be_bytes()); 1163 for target in targets { 1164 digest.update(u32::from(target.ordinal).to_be_bytes()); 1165 digest.update( 1166 u64::try_from(target.relay_id.len()) 1167 .map_err(|_| OperationError::Invariant)? 1168 .to_be_bytes(), 1169 ); 1170 digest.update(target.relay_id.as_bytes()); 1171 digest.update([u8::from(target.required)]); 1172 } 1173 let actual: [u8; 32] = digest.finalize().into(); 1174 if actual != outbox.target_set_sha256 { 1175 return Err(OperationError::Invariant); 1176 } 1177 Ok(()) 1178 } 1179 1180 fn validate_authority( 1181 outbox: OutboxRecord, 1182 authority: AuthorityBinding, 1183 ) -> Result<(), OperationError> { 1184 if outbox.authority_sha256 != authority.authority_sha256 1185 || outbox.target_set_sha256 != authority.target_set_sha256 1186 || outbox.target_count != authority.target_count 1187 || outbox.required_target_count != authority.required_target_count 1188 || outbox.max_attempts != authority.max_attempts 1189 || outbox.initial_backoff_ms != authority.initial_backoff_ms 1190 || outbox.maximum_backoff_ms != authority.maximum_backoff_ms 1191 || outbox.attempt_deadline_ms != authority.attempt_deadline_ms 1192 { 1193 return Err(OperationError::Invariant); 1194 } 1195 Ok(()) 1196 } 1197 1198 fn disposition( 1199 outbox: OutboxRecord, 1200 targets: &[TargetRecord], 1201 ) -> Result<OutboxDisposition, OperationError> { 1202 validate_targets(outbox, targets)?; 1203 let required = targets.iter().filter(|target| target.required); 1204 if required 1205 .clone() 1206 .all(|target| target.state == RhiPublicationTargetState::Accepted) 1207 { 1208 return Ok(OutboxDisposition { 1209 state: RhiPublicationOutboxState::Complete, 1210 next_attempt: None, 1211 }); 1212 } 1213 if required 1214 .clone() 1215 .any(|target| target_is_blocking(target, outbox.max_attempts)) 1216 { 1217 return Ok(OutboxDisposition { 1218 state: RhiPublicationOutboxState::Blocked, 1219 next_attempt: None, 1220 }); 1221 } 1222 let next_attempt = required 1223 .filter_map(|target| target.next_attempt) 1224 .min() 1225 .ok_or(OperationError::Invariant)?; 1226 Ok(OutboxDisposition { 1227 state: RhiPublicationOutboxState::Pending, 1228 next_attempt: Some(next_attempt), 1229 }) 1230 } 1231 1232 fn target_is_blocking(target: &TargetRecord, max_attempts: u16) -> bool { 1233 matches!( 1234 target.state, 1235 RhiPublicationTargetState::Rejected | RhiPublicationTargetState::AuthRequired 1236 ) || (retryable_state(target.state) 1237 && target.attempt_count >= max_attempts 1238 && target.next_attempt.is_none()) 1239 } 1240 1241 fn target_is_due( 1242 target: &TargetRecord, 1243 max_attempts: u16, 1244 now: RhiPublicationUnixMilliseconds, 1245 ) -> bool { 1246 retryable_state(target.state) 1247 && target.attempt_count < max_attempts 1248 && target.next_attempt.is_some_and(|next| next <= now) 1249 } 1250 1251 const fn retryable_state(state: RhiPublicationTargetState) -> bool { 1252 matches!( 1253 state, 1254 RhiPublicationTargetState::Pending 1255 | RhiPublicationTargetState::Failed 1256 | RhiPublicationTargetState::RateLimited 1257 | RhiPublicationTargetState::Unknown 1258 ) 1259 } 1260 1261 const fn retryable_outcome(outcome: RhiPublicationAttemptOutcome) -> bool { 1262 matches!( 1263 outcome, 1264 RhiPublicationAttemptOutcome::RateLimited 1265 | RhiPublicationAttemptOutcome::Failed 1266 | RhiPublicationAttemptOutcome::Unknown 1267 ) 1268 } 1269 1270 fn target_outcome_schedule( 1271 input: RecordInput, 1272 ) -> Result< 1273 ( 1274 RhiPublicationTargetState, 1275 Option<RhiPublicationUnixMilliseconds>, 1276 ), 1277 OperationError, 1278 > { 1279 let state = input.outcome.target_state(); 1280 if input.outcome == RhiPublicationAttemptOutcome::Submitted { 1281 return Err(OperationError::InvalidInput); 1282 } 1283 let next = if retryable_outcome(input.outcome) 1284 && input.attempt_number < input.lease.outbox.max_attempts 1285 { 1286 if input.retry_delay.get() > retry_upper_bound(input.lease.outbox, input.attempt_number) { 1287 return Err(OperationError::InvalidInput); 1288 } 1289 Some(RhiPublicationUnixMilliseconds( 1290 input 1291 .finished_at 1292 .get() 1293 .checked_add(input.retry_delay.get()) 1294 .filter(|value| *value <= MAX_UNIX_MILLISECONDS) 1295 .ok_or(OperationError::InvalidInput)?, 1296 )) 1297 } else { 1298 if input.retry_delay.get() != 0 { 1299 return Err(OperationError::InvalidInput); 1300 } 1301 None 1302 }; 1303 Ok((state, next)) 1304 } 1305 1306 fn retry_upper_bound(outbox: OutboxRecord, attempt_number: u16) -> u64 { 1307 let mut bound = outbox.initial_backoff_ms; 1308 for _ in 1..attempt_number { 1309 bound = bound.saturating_mul(2).min(outbox.maximum_backoff_ms); 1310 } 1311 bound.min(outbox.maximum_backoff_ms) 1312 } 1313 1314 fn decode_outbox(row: sqlx::sqlite::SqliteRow) -> Result<OutboxRecord, OperationError> { 1315 let id = RhiPublicationOutboxId::from_committed_bytes(blob::<32>(&row, "outbox_id")?); 1316 let state = decode_outbox_state(bounded_text(&row, "state", "state_bytes", 8)?.as_str())?; 1317 let target_count = bounded_u16(&row, "target_count", 1, 32)? as u8; 1318 let required_target_count = 1319 bounded_u16(&row, "required_target_count", 0, target_count.into())? as u8; 1320 let max_attempts = bounded_u16( 1321 &row, 1322 "max_attempts", 1323 1, 1324 RHI_PUBLICATION_ATTEMPT_NUMBER_MAXIMUM, 1325 )?; 1326 let initial_backoff_ms = bounded_u64(&row, "initial_backoff_ms", 1, 60_000)?; 1327 let maximum_backoff_ms = bounded_u64(&row, "maximum_backoff_ms", 1, 3_600_000)?; 1328 if initial_backoff_ms > maximum_backoff_ms { 1329 return Err(OperationError::Invariant); 1330 } 1331 let attempt_deadline_ms = bounded_u64(&row, "attempt_deadline_ms", 100, 30_000)?; 1332 let lease_owner = optional_owner(&row)?; 1333 let lease_expires = optional_millis(&row, "lease_expires_unix_ms")?; 1334 let next_attempt = optional_millis(&row, "next_attempt_unix_ms")?; 1335 if !valid_outbox_shape(state, next_attempt, lease_owner, lease_expires) { 1336 return Err(OperationError::Invariant); 1337 } 1338 Ok(OutboxRecord { 1339 id, 1340 event_sha256: blob::<32>(&row, "event_sha256")?, 1341 authority_sha256: blob::<32>(&row, "publication_authority_sha256")?, 1342 target_set_sha256: blob::<32>(&row, "target_set_sha256")?, 1343 target_count, 1344 required_target_count, 1345 max_attempts, 1346 initial_backoff_ms, 1347 maximum_backoff_ms, 1348 attempt_deadline_ms, 1349 state, 1350 revision: positive_u64(&row, "revision")?, 1351 next_attempt, 1352 lease_owner, 1353 lease_expires, 1354 created_at: nonnegative_millis(&row, "created_at_unix_ms")?, 1355 updated_at: nonnegative_millis(&row, "updated_at_unix_ms")?, 1356 }) 1357 } 1358 1359 fn decode_target(row: sqlx::sqlite::SqliteRow) -> Result<TargetRecord, OperationError> { 1360 let ordinal = bounded_u16( 1361 &row, 1362 "target_ordinal", 1363 0, 1364 u16::try_from(RHI_PUBLICATION_TARGET_ORDINAL_MAXIMUM) 1365 .map_err(|_| OperationError::Invariant)?, 1366 )? as u8; 1367 let relay_id = bounded_text(&row, "relay_id", "relay_id_bytes", 64)?; 1368 if relay_id.is_empty() 1369 || !relay_id.bytes().enumerate().all(|(index, byte)| { 1370 (byte.is_ascii_lowercase() || byte.is_ascii_digit() || matches!(byte, b'_' | b'-')) 1371 && (index != 0 || byte.is_ascii_lowercase()) 1372 }) 1373 { 1374 return Err(OperationError::Invariant); 1375 } 1376 let required = match row.try_get::<i64, _>("required") { 1377 Ok(0) => false, 1378 Ok(1) => true, 1379 _ => return Err(OperationError::Invariant), 1380 }; 1381 let state = decode_target_state(bounded_text(&row, "state", "state_bytes", 13)?.as_str())?; 1382 let attempt_count = bounded_u16( 1383 &row, 1384 "attempt_count", 1385 0, 1386 RHI_PUBLICATION_ATTEMPT_NUMBER_MAXIMUM, 1387 )?; 1388 let last_attempt_id = optional_attempt_id(&row)?; 1389 if (attempt_count == 0) != last_attempt_id.is_none() { 1390 return Err(OperationError::Invariant); 1391 } 1392 Ok(TargetRecord { 1393 ordinal, 1394 relay_id: relay_id.into_boxed_str(), 1395 required, 1396 state, 1397 revision: positive_u64(&row, "revision")?, 1398 attempt_count, 1399 next_attempt: optional_millis(&row, "next_attempt_unix_ms")?, 1400 last_attempt_id, 1401 updated_at: nonnegative_millis(&row, "updated_at_unix_ms")?, 1402 }) 1403 } 1404 1405 async fn read_attempt( 1406 transaction: &mut ServiceSqliteTransaction<'_>, 1407 id: RhiPublicationAttemptId, 1408 ) -> Result<Option<AttemptRecord>, OperationError> { 1409 let rows = sqlx::query(READ_ATTEMPT_SQL) 1410 .bind(id.as_bytes().as_slice()) 1411 .fetch_all(&mut *transaction) 1412 .await 1413 .map_err(|_| OperationError::Storage)?; 1414 exactly_zero_or_one(rows)?.map(decode_attempt).transpose() 1415 } 1416 1417 fn decode_attempt(row: sqlx::sqlite::SqliteRow) -> Result<AttemptRecord, OperationError> { 1418 let outcome = 1419 decode_attempt_outcome(bounded_text(&row, "outcome", "outcome_bytes", 13)?.as_str())?; 1420 let result_code = bounded_text(&row, "result_code", "result_code_bytes", 64)?; 1421 if result_code != outcome.code() { 1422 return Err(OperationError::Invariant); 1423 } 1424 Ok(AttemptRecord { 1425 attempt_id: attempt_id_from_durable_bytes(blob::<32>(&row, "attempt_id")?), 1426 outbox_id: RhiPublicationOutboxId::from_committed_bytes(blob::<32>(&row, "outbox_id")?), 1427 target_ordinal: bounded_u16(&row, "target_ordinal", 0, 31)? as u8, 1428 attempt_number: bounded_u16(&row, "attempt_number", 1, 100)?, 1429 event_sha256: blob::<32>(&row, "event_sha256")?, 1430 lease_owner: RhiPublicationLeaseOwner(blob::<16>(&row, "lease_owner")?), 1431 started_at: nonnegative_millis(&row, "started_at_unix_ms")?, 1432 finished_at: nonnegative_millis(&row, "finished_at_unix_ms")?, 1433 outcome, 1434 }) 1435 } 1436 1437 fn decode_outbox_state(value: &str) -> Result<RhiPublicationOutboxState, OperationError> { 1438 match value { 1439 "pending" => Ok(RhiPublicationOutboxState::Pending), 1440 "leased" => Ok(RhiPublicationOutboxState::Leased), 1441 "complete" => Ok(RhiPublicationOutboxState::Complete), 1442 "blocked" => Ok(RhiPublicationOutboxState::Blocked), 1443 _ => Err(OperationError::Invariant), 1444 } 1445 } 1446 1447 fn decode_target_state(value: &str) -> Result<RhiPublicationTargetState, OperationError> { 1448 match value { 1449 "pending" => Ok(RhiPublicationTargetState::Pending), 1450 "submitted" => Ok(RhiPublicationTargetState::Submitted), 1451 "accepted" => Ok(RhiPublicationTargetState::Accepted), 1452 "rejected" => Ok(RhiPublicationTargetState::Rejected), 1453 "rate_limited" => Ok(RhiPublicationTargetState::RateLimited), 1454 "auth_required" => Ok(RhiPublicationTargetState::AuthRequired), 1455 "failed" => Ok(RhiPublicationTargetState::Failed), 1456 "unknown" => Ok(RhiPublicationTargetState::Unknown), 1457 _ => Err(OperationError::Invariant), 1458 } 1459 } 1460 1461 fn decode_attempt_outcome(value: &str) -> Result<RhiPublicationAttemptOutcome, OperationError> { 1462 match value { 1463 "submitted" => Ok(RhiPublicationAttemptOutcome::Submitted), 1464 "accepted" => Ok(RhiPublicationAttemptOutcome::Accepted), 1465 "rejected" => Ok(RhiPublicationAttemptOutcome::Rejected), 1466 "rate_limited" => Ok(RhiPublicationAttemptOutcome::RateLimited), 1467 "auth_required" => Ok(RhiPublicationAttemptOutcome::AuthRequired), 1468 "failed" => Ok(RhiPublicationAttemptOutcome::Failed), 1469 "unknown" => Ok(RhiPublicationAttemptOutcome::Unknown), 1470 _ => Err(OperationError::Invariant), 1471 } 1472 } 1473 1474 fn valid_outbox_shape( 1475 state: RhiPublicationOutboxState, 1476 next: Option<RhiPublicationUnixMilliseconds>, 1477 owner: Option<RhiPublicationLeaseOwner>, 1478 expires: Option<RhiPublicationUnixMilliseconds>, 1479 ) -> bool { 1480 match state { 1481 RhiPublicationOutboxState::Pending => { 1482 next.is_some() && owner.is_none() && expires.is_none() 1483 } 1484 RhiPublicationOutboxState::Leased => { 1485 next.is_none() && owner.is_some() && expires.is_some_and(|value| value.get() > 0) 1486 } 1487 RhiPublicationOutboxState::Complete | RhiPublicationOutboxState::Blocked => { 1488 next.is_none() && owner.is_none() && expires.is_none() 1489 } 1490 } 1491 } 1492 1493 fn exactly_zero_or_one( 1494 mut rows: Vec<sqlx::sqlite::SqliteRow>, 1495 ) -> Result<Option<sqlx::sqlite::SqliteRow>, OperationError> { 1496 match rows.len() { 1497 0 => Ok(None), 1498 1 => Ok(rows.pop()), 1499 _ => Err(OperationError::Invariant), 1500 } 1501 } 1502 1503 fn blob<const N: usize>( 1504 row: &sqlx::sqlite::SqliteRow, 1505 column: &str, 1506 ) -> Result<[u8; N], OperationError> { 1507 row.try_get::<Option<Vec<u8>>, _>(column) 1508 .map_err(|_| OperationError::Invariant)? 1509 .ok_or(OperationError::Invariant)? 1510 .try_into() 1511 .map_err(|_| OperationError::Invariant) 1512 } 1513 1514 fn bounded_text( 1515 row: &sqlx::sqlite::SqliteRow, 1516 column: &str, 1517 length_column: &str, 1518 maximum: usize, 1519 ) -> Result<String, OperationError> { 1520 let length = row 1521 .try_get::<i64, _>(length_column) 1522 .ok() 1523 .and_then(|value| usize::try_from(value).ok()) 1524 .filter(|value| *value <= maximum) 1525 .ok_or(OperationError::Invariant)?; 1526 let value = row 1527 .try_get::<String, _>(column) 1528 .map_err(|_| OperationError::Invariant)?; 1529 if value.len() != length { 1530 return Err(OperationError::Invariant); 1531 } 1532 Ok(value) 1533 } 1534 1535 fn bounded_u16( 1536 row: &sqlx::sqlite::SqliteRow, 1537 column: &str, 1538 minimum: u16, 1539 maximum: u16, 1540 ) -> Result<u16, OperationError> { 1541 row.try_get::<i64, _>(column) 1542 .ok() 1543 .and_then(|value| u16::try_from(value).ok()) 1544 .filter(|value| *value >= minimum && *value <= maximum) 1545 .ok_or(OperationError::Invariant) 1546 } 1547 1548 fn bounded_u64( 1549 row: &sqlx::sqlite::SqliteRow, 1550 column: &str, 1551 minimum: u64, 1552 maximum: u64, 1553 ) -> Result<u64, OperationError> { 1554 row.try_get::<i64, _>(column) 1555 .ok() 1556 .and_then(|value| u64::try_from(value).ok()) 1557 .filter(|value| *value >= minimum && *value <= maximum) 1558 .ok_or(OperationError::Invariant) 1559 } 1560 1561 fn positive_u64(row: &sqlx::sqlite::SqliteRow, column: &str) -> Result<u64, OperationError> { 1562 bounded_u64(row, column, 1, MAX_UNIX_MILLISECONDS) 1563 } 1564 1565 fn nonnegative_millis( 1566 row: &sqlx::sqlite::SqliteRow, 1567 column: &str, 1568 ) -> Result<RhiPublicationUnixMilliseconds, OperationError> { 1569 bounded_u64(row, column, 0, MAX_UNIX_MILLISECONDS).map(RhiPublicationUnixMilliseconds) 1570 } 1571 1572 fn optional_millis( 1573 row: &sqlx::sqlite::SqliteRow, 1574 column: &str, 1575 ) -> Result<Option<RhiPublicationUnixMilliseconds>, OperationError> { 1576 row.try_get::<Option<i64>, _>(column) 1577 .map_err(|_| OperationError::Invariant)? 1578 .map(|value| { 1579 u64::try_from(value) 1580 .ok() 1581 .filter(|value| *value <= MAX_UNIX_MILLISECONDS) 1582 .map(RhiPublicationUnixMilliseconds) 1583 .ok_or(OperationError::Invariant) 1584 }) 1585 .transpose() 1586 } 1587 1588 fn optional_owner( 1589 row: &sqlx::sqlite::SqliteRow, 1590 ) -> Result<Option<RhiPublicationLeaseOwner>, OperationError> { 1591 let length = row 1592 .try_get::<Option<i64>, _>("lease_owner_bytes") 1593 .map_err(|_| OperationError::Invariant)?; 1594 let bytes = row 1595 .try_get::<Option<Vec<u8>>, _>("lease_owner") 1596 .map_err(|_| OperationError::Invariant)?; 1597 match (length, bytes) { 1598 (None, None) => Ok(None), 1599 (Some(length), Some(bytes)) 1600 if length == LEASE_OWNER_BYTES as i64 && bytes.len() == LEASE_OWNER_BYTES => 1601 { 1602 let bytes: [u8; LEASE_OWNER_BYTES] = 1603 bytes.try_into().map_err(|_| OperationError::Invariant)?; 1604 RhiPublicationLeaseOwner::from_bytes(bytes) 1605 .map(Some) 1606 .map_err(|_| OperationError::Invariant) 1607 } 1608 _ => Err(OperationError::Invariant), 1609 } 1610 } 1611 1612 fn optional_attempt_id( 1613 row: &sqlx::sqlite::SqliteRow, 1614 ) -> Result<Option<RhiPublicationAttemptId>, OperationError> { 1615 let length = row 1616 .try_get::<Option<i64>, _>("last_attempt_id_bytes") 1617 .map_err(|_| OperationError::Invariant)?; 1618 let bytes = row 1619 .try_get::<Option<Vec<u8>>, _>("last_attempt_id") 1620 .map_err(|_| OperationError::Invariant)?; 1621 match (length, bytes) { 1622 (None, None) => Ok(None), 1623 (Some(32), Some(bytes)) if bytes.len() == 32 => Ok(Some(attempt_id_from_durable_bytes( 1624 bytes.try_into().map_err(|_| OperationError::Invariant)?, 1625 ))), 1626 _ => Err(OperationError::Invariant), 1627 } 1628 } 1629 1630 fn attempt_id_from_durable_bytes(bytes: [u8; 32]) -> RhiPublicationAttemptId { 1631 // Restricted to bounded database decoding; every caller subsequently 1632 // compares the value with a freshly domain-derived identity. 1633 crate::publication_attempt::attempt_id_from_durable_bytes(bytes) 1634 } 1635 1636 fn i64_value(value: u64) -> Result<i64, OperationError> { 1637 i64::try_from(value).map_err(|_| OperationError::InvalidInput) 1638 } 1639 1640 fn optional_i64( 1641 value: Option<RhiPublicationUnixMilliseconds>, 1642 ) -> Result<Option<i64>, OperationError> { 1643 value.map(|value| i64_value(value.get())).transpose() 1644 } 1645 1646 #[derive(Clone, Copy, PartialEq, Eq)] 1647 struct OutboxRecord { 1648 id: RhiPublicationOutboxId, 1649 event_sha256: [u8; 32], 1650 authority_sha256: [u8; 32], 1651 target_set_sha256: [u8; 32], 1652 target_count: u8, 1653 required_target_count: u8, 1654 max_attempts: u16, 1655 initial_backoff_ms: u64, 1656 maximum_backoff_ms: u64, 1657 attempt_deadline_ms: u64, 1658 state: RhiPublicationOutboxState, 1659 revision: u64, 1660 next_attempt: Option<RhiPublicationUnixMilliseconds>, 1661 lease_owner: Option<RhiPublicationLeaseOwner>, 1662 lease_expires: Option<RhiPublicationUnixMilliseconds>, 1663 created_at: RhiPublicationUnixMilliseconds, 1664 updated_at: RhiPublicationUnixMilliseconds, 1665 } 1666 1667 #[derive(Clone, Copy)] 1668 struct AuthorityBinding { 1669 authority_sha256: [u8; 32], 1670 target_set_sha256: [u8; 32], 1671 target_count: u8, 1672 required_target_count: u8, 1673 max_attempts: u16, 1674 initial_backoff_ms: u64, 1675 maximum_backoff_ms: u64, 1676 attempt_deadline_ms: u64, 1677 } 1678 1679 impl AuthorityBinding { 1680 fn from_authority(authority: &RhiPublicationAuthority) -> Option<Self> { 1681 if authority.mode() == RhiPublicationMode::Disabled { 1682 return None; 1683 } 1684 let retry = authority.retry_policy()?; 1685 Some(Self { 1686 authority_sha256: *authority.authority_sha256(), 1687 target_set_sha256: *authority.target_set_sha256(), 1688 target_count: u8::try_from(authority.targets().len()).ok()?, 1689 required_target_count: u8::try_from( 1690 authority 1691 .targets() 1692 .iter() 1693 .filter(|target| target.required()) 1694 .count(), 1695 ) 1696 .ok()?, 1697 max_attempts: retry.maximum_attempts(), 1698 initial_backoff_ms: retry.initial_backoff_milliseconds(), 1699 maximum_backoff_ms: retry.maximum_backoff_milliseconds(), 1700 attempt_deadline_ms: retry.attempt_deadline_milliseconds(), 1701 }) 1702 } 1703 } 1704 1705 #[derive(PartialEq, Eq)] 1706 struct TargetRecord { 1707 ordinal: u8, 1708 relay_id: Box<str>, 1709 required: bool, 1710 state: RhiPublicationTargetState, 1711 revision: u64, 1712 attempt_count: u16, 1713 next_attempt: Option<RhiPublicationUnixMilliseconds>, 1714 last_attempt_id: Option<RhiPublicationAttemptId>, 1715 updated_at: RhiPublicationUnixMilliseconds, 1716 } 1717 1718 impl Clone for TargetRecord { 1719 fn clone(&self) -> Self { 1720 Self { 1721 ordinal: self.ordinal, 1722 relay_id: self.relay_id.clone(), 1723 required: self.required, 1724 state: self.state, 1725 revision: self.revision, 1726 attempt_count: self.attempt_count, 1727 next_attempt: self.next_attempt, 1728 last_attempt_id: self.last_attempt_id, 1729 updated_at: self.updated_at, 1730 } 1731 } 1732 } 1733 1734 #[derive(Clone, Copy)] 1735 struct RecoveryCandidate { 1736 outbox: OutboxRecord, 1737 retry_upper_bound: u64, 1738 } 1739 1740 #[derive(Clone, Copy)] 1741 struct RecordInput { 1742 lease: RhiPublicationLease, 1743 target_ordinal: u8, 1744 target_revision: u64, 1745 attempt_number: u16, 1746 attempt_id: RhiPublicationAttemptId, 1747 started_at: RhiPublicationUnixMilliseconds, 1748 finished_at: RhiPublicationUnixMilliseconds, 1749 outcome: RhiPublicationAttemptOutcome, 1750 retry_delay: RhiPublicationRetryDelayMilliseconds, 1751 } 1752 1753 impl RecordInput { 1754 fn from_prepared( 1755 prepared: &RhiPreparedPublicationAttempt, 1756 finished_at: RhiPublicationUnixMilliseconds, 1757 outcome: RhiPublicationAttemptOutcome, 1758 retry_delay: RhiPublicationRetryDelayMilliseconds, 1759 ) -> Result<Self, RhiPublicationExecutionError> { 1760 if finished_at < prepared.started_at || outcome == RhiPublicationAttemptOutcome::Submitted { 1761 return Err(RhiPublicationExecutionError::new( 1762 RhiPublicationExecutionErrorKind::InvalidInput, 1763 )); 1764 } 1765 let evidence = RhiPublicationAttemptEvidence::new( 1766 &prepared.publication, 1767 u32::from(prepared.target.ordinal), 1768 prepared.target.attempt_count, 1769 prepared.started_at, 1770 finished_at, 1771 outcome, 1772 ) 1773 .map_err(|_| { 1774 RhiPublicationExecutionError::new(RhiPublicationExecutionErrorKind::InvalidInput) 1775 })?; 1776 if evidence.id() != prepared.attempt_id { 1777 return Err(RhiPublicationExecutionError::new( 1778 RhiPublicationExecutionErrorKind::Invariant, 1779 )); 1780 } 1781 Ok(Self { 1782 lease: prepared.lease, 1783 target_ordinal: prepared.target.ordinal, 1784 target_revision: prepared.target.revision, 1785 attempt_number: prepared.target.attempt_count, 1786 attempt_id: prepared.attempt_id, 1787 started_at: prepared.started_at, 1788 finished_at, 1789 outcome, 1790 retry_delay, 1791 }) 1792 } 1793 } 1794 1795 #[derive(Clone, Copy, PartialEq, Eq)] 1796 struct AttemptRecord { 1797 attempt_id: RhiPublicationAttemptId, 1798 outbox_id: RhiPublicationOutboxId, 1799 target_ordinal: u8, 1800 attempt_number: u16, 1801 event_sha256: [u8; 32], 1802 lease_owner: RhiPublicationLeaseOwner, 1803 started_at: RhiPublicationUnixMilliseconds, 1804 finished_at: RhiPublicationUnixMilliseconds, 1805 outcome: RhiPublicationAttemptOutcome, 1806 } 1807 1808 impl AttemptRecord { 1809 const fn from_input(input: RecordInput) -> Self { 1810 Self { 1811 attempt_id: input.attempt_id, 1812 outbox_id: input.lease.outbox.id, 1813 target_ordinal: input.target_ordinal, 1814 attempt_number: input.attempt_number, 1815 event_sha256: input.lease.outbox.event_sha256, 1816 lease_owner: input.lease.owner, 1817 started_at: input.started_at, 1818 finished_at: input.finished_at, 1819 outcome: input.outcome, 1820 } 1821 } 1822 } 1823 1824 #[derive(Clone, Copy)] 1825 struct OutboxDisposition { 1826 state: RhiPublicationOutboxState, 1827 next_attempt: Option<RhiPublicationUnixMilliseconds>, 1828 } 1829 1830 #[derive(Clone, Copy, Debug, PartialEq, Eq)] 1831 enum OperationError { 1832 InvalidInput, 1833 NotReady, 1834 LeaseLost, 1835 Invariant, 1836 Storage, 1837 } 1838 1839 fn map_transaction_error( 1840 error: ServiceSqliteTransactionError<OperationError>, 1841 ) -> RhiPublicationExecutionError { 1842 if error.kind() == ServiceSqliteTransactionErrorKind::CommitOutcomeUnknown { 1843 return RhiPublicationExecutionError::new( 1844 RhiPublicationExecutionErrorKind::CommitOutcomeUnknown, 1845 ); 1846 } 1847 RhiPublicationExecutionError::new(match error.operation_error().copied() { 1848 Some(OperationError::InvalidInput) => RhiPublicationExecutionErrorKind::InvalidInput, 1849 Some(OperationError::NotReady) => RhiPublicationExecutionErrorKind::NotReady, 1850 Some(OperationError::LeaseLost) => RhiPublicationExecutionErrorKind::LeaseLost, 1851 Some(OperationError::Invariant) => RhiPublicationExecutionErrorKind::Invariant, 1852 Some(OperationError::Storage) | None => RhiPublicationExecutionErrorKind::Storage, 1853 }) 1854 } 1855 1856 #[cfg(test)] 1857 mod tests { 1858 use super::*; 1859 1860 #[test] 1861 fn retry_bounds_are_saturating_closed_and_errors_are_safe() { 1862 let outbox = OutboxRecord { 1863 id: RhiPublicationOutboxId::from_committed_bytes([1; 32]), 1864 event_sha256: [2; 32], 1865 authority_sha256: [3; 32], 1866 target_set_sha256: [4; 32], 1867 target_count: 1, 1868 required_target_count: 1, 1869 max_attempts: 100, 1870 initial_backoff_ms: 250, 1871 maximum_backoff_ms: 30_000, 1872 attempt_deadline_ms: 5_000, 1873 state: RhiPublicationOutboxState::Pending, 1874 revision: 1, 1875 next_attempt: Some(RhiPublicationUnixMilliseconds(0)), 1876 lease_owner: None, 1877 lease_expires: None, 1878 created_at: RhiPublicationUnixMilliseconds(0), 1879 updated_at: RhiPublicationUnixMilliseconds(0), 1880 }; 1881 assert_eq!(retry_upper_bound(outbox, 1), 250); 1882 assert_eq!(retry_upper_bound(outbox, 2), 500); 1883 assert_eq!(retry_upper_bound(outbox, 100), 30_000); 1884 assert!(RhiPublicationLeaseOwner::from_bytes([0; 16]).is_err()); 1885 assert!(RhiPublicationRetryDelayMilliseconds::new(3_600_001).is_err()); 1886 for kind in [ 1887 RhiPublicationExecutionErrorKind::InvalidMode, 1888 RhiPublicationExecutionErrorKind::InvalidInput, 1889 RhiPublicationExecutionErrorKind::NotReady, 1890 RhiPublicationExecutionErrorKind::LeaseLost, 1891 RhiPublicationExecutionErrorKind::Invariant, 1892 RhiPublicationExecutionErrorKind::ClockUnavailable, 1893 RhiPublicationExecutionErrorKind::EntropyUnavailable, 1894 RhiPublicationExecutionErrorKind::Storage, 1895 RhiPublicationExecutionErrorKind::CommitOutcomeUnknown, 1896 ] { 1897 let error = RhiPublicationExecutionError::new(kind); 1898 assert!(error.code().starts_with("publication_execution_")); 1899 assert!(Error::source(&error).is_none()); 1900 let rendered = format!("{error} {error:?}"); 1901 assert!(!rendered.contains("relay-primary")); 1902 assert!(!rendered.contains("SELECT")); 1903 } 1904 } 1905 }