state_delivery.rs (79975B)
1 //! Durable delivery-job identity, bounded target attempts, and lease evidence. 2 3 use core::fmt; 4 use std::error::Error; 5 6 use radroots_service_sqlite::{ 7 ServiceSqliteTransaction, ServiceSqliteTransactionError, ServiceSqliteTransactionErrorKind, 8 }; 9 use sha2::{Digest, Sha256}; 10 use sqlx::Row; 11 12 use crate::MycSignerOperationId; 13 use crate::state_repository::{ 14 MycStateRepository, MycStateRepositoryError, MycStateRepositoryErrorKind, PersistedMetadata, 15 RepositoryOperationError, require_expected_metadata, 16 }; 17 18 /// Maximum immutable relay targets on one delivery job. 19 pub const MYC_DELIVERY_TARGET_MAX_COUNT: usize = 32; 20 /// Maximum durable attempts for one target. 21 pub const MYC_DELIVERY_ATTEMPT_MAX_COUNT: u32 = 32; 22 /// Maximum UTF-8 byte length of a canonical delivery relay identifier. 23 pub const MYC_DELIVERY_RELAY_ID_MAX_BYTES: usize = 64; 24 /// Maximum caller-injected full-jitter delay for one retry. 25 pub const MYC_DELIVERY_RETRY_JITTER_MAX_MS: u64 = 300_000; 26 27 const JOB_ID_DOMAIN: &[u8] = b"radroots.myc.delivery_job.v1\0"; 28 const ATTEMPT_ID_DOMAIN: &[u8] = b"radroots.myc.delivery_attempt.v1\0"; 29 30 const READ_ACTIVE_JOB_COUNT_SQL: &str = 31 "SELECT COUNT(*) AS row_count FROM delivery_jobs WHERE status IN ('pending', 'active')"; 32 33 const READ_RUNTIME_OUTBOX_STATUS_SQL: &str = r#"SELECT 34 COUNT(CASE WHEN status IN ('pending', 'active') THEN 1 END) AS pending_count, 35 (SELECT COUNT(*) FROM delivery_targets WHERE status = 'unknown') AS unknown_count, 36 MIN(CASE WHEN status IN ('pending', 'active') THEN created_at_unix_ms ELSE NULL END) 37 AS oldest_pending_at_unix_ms, 38 typeof(MIN(CASE WHEN status IN ('pending', 'active') THEN created_at_unix_ms ELSE NULL END)) 39 AS oldest_pending_type 40 FROM delivery_jobs"#; 41 42 const READ_NEXT_READY_TARGET_SQL: &str = r#"SELECT 43 CASE WHEN typeof(t.job_id) = 'blob' AND length(t.job_id) = 32 44 THEN t.job_id ELSE NULL END AS job_id, 45 CASE WHEN typeof(t.relay_id) = 'text' 46 AND length(CAST(t.relay_id AS BLOB)) BETWEEN 1 AND 64 47 THEN t.relay_id ELSE NULL END AS relay_id 48 FROM delivery_targets t 49 JOIN delivery_jobs j ON j.job_id = t.job_id 50 WHERE j.status IN ('pending', 'active') 51 AND t.status IN ('pending', 'retryable', 'unknown') 52 AND t.active_attempt_id IS NULL 53 AND (t.next_attempt_at_unix_ms IS NULL OR t.next_attempt_at_unix_ms <= ?) 54 ORDER BY j.created_at_unix_ms, t.job_id, t.target_index 55 LIMIT 2"#; 56 57 const READ_JOB_SQL: &str = r#"SELECT 58 CASE WHEN typeof(job_id) = 'blob' AND length(job_id) = 32 59 THEN job_id ELSE NULL END AS job_id, 60 CASE WHEN typeof(source_kind) = 'text' AND length(CAST(source_kind AS BLOB)) <= 32 61 THEN source_kind ELSE NULL END AS source_kind, 62 CASE WHEN typeof(source_id) = 'blob' AND length(source_id) = 32 63 THEN source_id ELSE NULL END AS source_id, 64 CASE WHEN typeof(artifact_sha256) = 'blob' AND length(artifact_sha256) = 32 65 THEN artifact_sha256 ELSE NULL END AS artifact_sha256, 66 CASE WHEN typeof(policy_mode) = 'text' AND length(CAST(policy_mode AS BLOB)) <= 32 67 THEN policy_mode ELSE NULL END AS policy_mode, 68 required_acknowledgements, max_attempts, initial_backoff_ms, 69 maximum_backoff_ms, attempt_deadline_ms, 70 CASE WHEN typeof(status) = 'text' AND length(CAST(status AS BLOB)) <= 16 71 THEN status ELSE NULL END AS status, 72 created_at_unix_ms, updated_at_unix_ms, finalized_at_unix_ms, 73 typeof(finalized_at_unix_ms) AS finalized_at_type 74 FROM delivery_jobs 75 WHERE job_id = ? 76 LIMIT 2"#; 77 78 const READ_JOB_BY_SOURCE_SQL: &str = r#"SELECT 79 CASE WHEN typeof(job_id) = 'blob' AND length(job_id) = 32 80 THEN job_id ELSE NULL END AS job_id, 81 CASE WHEN typeof(source_kind) = 'text' AND length(CAST(source_kind AS BLOB)) <= 32 82 THEN source_kind ELSE NULL END AS source_kind, 83 CASE WHEN typeof(source_id) = 'blob' AND length(source_id) = 32 84 THEN source_id ELSE NULL END AS source_id, 85 CASE WHEN typeof(artifact_sha256) = 'blob' AND length(artifact_sha256) = 32 86 THEN artifact_sha256 ELSE NULL END AS artifact_sha256, 87 CASE WHEN typeof(policy_mode) = 'text' AND length(CAST(policy_mode AS BLOB)) <= 32 88 THEN policy_mode ELSE NULL END AS policy_mode, 89 required_acknowledgements, max_attempts, initial_backoff_ms, 90 maximum_backoff_ms, attempt_deadline_ms, 91 CASE WHEN typeof(status) = 'text' AND length(CAST(status AS BLOB)) <= 16 92 THEN status ELSE NULL END AS status, 93 created_at_unix_ms, updated_at_unix_ms, finalized_at_unix_ms, 94 typeof(finalized_at_unix_ms) AS finalized_at_type 95 FROM delivery_jobs 96 WHERE source_kind = ? AND source_id = ? 97 LIMIT 2"#; 98 99 const READ_TARGETS_SQL: &str = r#"SELECT 100 target_index, 101 CASE WHEN typeof(relay_id) = 'text' 102 AND length(CAST(relay_id AS BLOB)) BETWEEN 1 AND 64 103 THEN relay_id ELSE NULL END AS relay_id, 104 required, attempt_count, 105 CASE WHEN typeof(status) = 'text' AND length(CAST(status AS BLOB)) <= 16 106 THEN status ELSE NULL END AS status, 107 CASE WHEN typeof(active_attempt_id) = 'blob' AND length(active_attempt_id) = 32 108 THEN active_attempt_id ELSE NULL END AS active_attempt_id, 109 typeof(active_attempt_id) AS active_attempt_id_type, 110 next_attempt_at_unix_ms, 111 typeof(next_attempt_at_unix_ms) AS next_attempt_at_type, 112 updated_at_unix_ms 113 FROM delivery_targets 114 WHERE job_id = ? 115 ORDER BY target_index 116 LIMIT 33"#; 117 118 const READ_ATTEMPTS_SQL: &str = r#"SELECT 119 CASE WHEN typeof(attempt_id) = 'blob' AND length(attempt_id) = 32 120 THEN attempt_id ELSE NULL END AS attempt_id, 121 attempt_number, 122 CASE WHEN typeof(attempt_nonce) = 'blob' AND length(attempt_nonce) = 32 123 THEN attempt_nonce ELSE NULL END AS attempt_nonce, 124 CASE WHEN typeof(status) = 'text' AND length(CAST(status AS BLOB)) <= 16 125 THEN status ELSE NULL END AS status, 126 leased_at_unix_ms, lease_expires_at_unix_ms, 127 submitted_at_unix_ms, typeof(submitted_at_unix_ms) AS submitted_at_type, 128 resolved_at_unix_ms, typeof(resolved_at_unix_ms) AS resolved_at_type, 129 CASE WHEN typeof(reason_code) = 'text' AND length(CAST(reason_code AS BLOB)) <= 32 130 THEN reason_code ELSE NULL END AS reason_code, 131 typeof(reason_code) AS reason_code_type 132 FROM delivery_attempts 133 WHERE job_id = ? AND target_index = ? 134 ORDER BY attempt_number 135 LIMIT 33"#; 136 137 const READ_ATTEMPT_SQL: &str = r#"SELECT 138 CASE WHEN typeof(attempt_id) = 'blob' AND length(attempt_id) = 32 139 THEN attempt_id ELSE NULL END AS attempt_id, 140 attempt_number, 141 CASE WHEN typeof(attempt_nonce) = 'blob' AND length(attempt_nonce) = 32 142 THEN attempt_nonce ELSE NULL END AS attempt_nonce, 143 CASE WHEN typeof(status) = 'text' AND length(CAST(status AS BLOB)) <= 16 144 THEN status ELSE NULL END AS status, 145 leased_at_unix_ms, lease_expires_at_unix_ms, 146 submitted_at_unix_ms, typeof(submitted_at_unix_ms) AS submitted_at_type, 147 resolved_at_unix_ms, typeof(resolved_at_unix_ms) AS resolved_at_type, 148 CASE WHEN typeof(reason_code) = 'text' AND length(CAST(reason_code AS BLOB)) <= 32 149 THEN reason_code ELSE NULL END AS reason_code, 150 typeof(reason_code) AS reason_code_type 151 FROM delivery_attempts 152 WHERE job_id = ? AND target_index = ? AND attempt_id = ? 153 LIMIT 2"#; 154 155 const READ_ATTEMPT_BY_NONCE_SQL: &str = r#"SELECT 156 CASE WHEN typeof(attempt_id) = 'blob' AND length(attempt_id) = 32 157 THEN attempt_id ELSE NULL END AS attempt_id, 158 attempt_number, 159 CASE WHEN typeof(attempt_nonce) = 'blob' AND length(attempt_nonce) = 32 160 THEN attempt_nonce ELSE NULL END AS attempt_nonce, 161 CASE WHEN typeof(status) = 'text' AND length(CAST(status AS BLOB)) <= 16 162 THEN status ELSE NULL END AS status, 163 leased_at_unix_ms, lease_expires_at_unix_ms, 164 submitted_at_unix_ms, typeof(submitted_at_unix_ms) AS submitted_at_type, 165 resolved_at_unix_ms, typeof(resolved_at_unix_ms) AS resolved_at_type, 166 CASE WHEN typeof(reason_code) = 'text' AND length(CAST(reason_code AS BLOB)) <= 32 167 THEN reason_code ELSE NULL END AS reason_code, 168 typeof(reason_code) AS reason_code_type 169 FROM delivery_attempts 170 WHERE job_id = ? AND target_index = ? AND attempt_nonce = ? 171 LIMIT 2"#; 172 173 const INSERT_JOB_SQL: &str = r#"INSERT INTO delivery_jobs ( 174 job_id, source_kind, source_id, artifact_sha256, policy_mode, 175 required_acknowledgements, max_attempts, initial_backoff_ms, 176 maximum_backoff_ms, attempt_deadline_ms, status, 177 created_at_unix_ms, updated_at_unix_ms, finalized_at_unix_ms 178 ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, 'pending', ?, ?, NULL)"#; 179 180 const INSERT_TARGET_SQL: &str = r#"INSERT INTO delivery_targets ( 181 job_id, target_index, relay_id, required, attempt_count, status, 182 active_attempt_id, next_attempt_at_unix_ms, updated_at_unix_ms 183 ) VALUES (?, ?, ?, ?, 0, 'pending', NULL, NULL, ?)"#; 184 185 const INSERT_ATTEMPT_SQL: &str = r#"INSERT INTO delivery_attempts ( 186 attempt_id, job_id, target_index, attempt_number, attempt_nonce, status, 187 leased_at_unix_ms, lease_expires_at_unix_ms, submitted_at_unix_ms, 188 resolved_at_unix_ms, reason_code 189 ) VALUES (?, ?, ?, ?, ?, 'leased', ?, ?, NULL, NULL, NULL)"#; 190 191 const CLAIM_TARGET_SQL: &str = r#"UPDATE delivery_targets 192 SET attempt_count = ?, status = 'leased', active_attempt_id = ?, 193 next_attempt_at_unix_ms = NULL, updated_at_unix_ms = ? 194 WHERE job_id = ? AND target_index = ? AND attempt_count = ? 195 AND status IN ('pending', 'retryable', 'unknown') 196 AND active_attempt_id IS NULL"#; 197 198 const MARK_JOB_ACTIVE_SQL: &str = r#"UPDATE delivery_jobs 199 SET status = 'active', updated_at_unix_ms = ? 200 WHERE job_id = ? AND status = 'pending'"#; 201 202 const MARK_ATTEMPT_SUBMITTED_SQL: &str = r#"UPDATE delivery_attempts 203 SET status = 'submitted', submitted_at_unix_ms = ? 204 WHERE attempt_id = ? AND job_id = ? AND target_index = ? 205 AND status = 'leased' AND lease_expires_at_unix_ms >= ?"#; 206 207 const MARK_TARGET_SUBMITTED_SQL: &str = r#"UPDATE delivery_targets 208 SET status = 'submitted', updated_at_unix_ms = ? 209 WHERE job_id = ? AND target_index = ? AND active_attempt_id = ? AND status = 'leased'"#; 210 211 const RESOLVE_ATTEMPT_SQL: &str = r#"UPDATE delivery_attempts 212 SET status = ?, resolved_at_unix_ms = ?, reason_code = ? 213 WHERE attempt_id = ? AND job_id = ? AND target_index = ? AND status = ?"#; 214 215 const RESOLVE_TARGET_SQL: &str = r#"UPDATE delivery_targets 216 SET status = ?, active_attempt_id = NULL, next_attempt_at_unix_ms = ?, 217 updated_at_unix_ms = ? 218 WHERE job_id = ? AND target_index = ? AND active_attempt_id = ? AND status = ?"#; 219 220 const FINALIZE_JOB_SQL: &str = r#"UPDATE delivery_jobs 221 SET status = ?, updated_at_unix_ms = ?, finalized_at_unix_ms = ? 222 WHERE job_id = ? AND status IN ('pending', 'active')"#; 223 224 /// Stable source-free construction failure classes. 225 #[derive(Clone, Copy, Debug, PartialEq, Eq)] 226 pub enum MycDeliveryStateErrorKind { 227 InvalidTime, 228 InvalidRelayId, 229 InvalidPolicy, 230 InvalidRetryJitter, 231 } 232 233 impl MycDeliveryStateErrorKind { 234 /// Returns the stable machine-readable classification. 235 #[must_use] 236 pub const fn code(self) -> &'static str { 237 match self { 238 Self::InvalidTime => "delivery_time_invalid", 239 Self::InvalidRelayId => "delivery_relay_id_invalid", 240 Self::InvalidPolicy => "delivery_policy_invalid", 241 Self::InvalidRetryJitter => "delivery_retry_jitter_invalid", 242 } 243 } 244 } 245 246 /// Source-free delivery-state construction failure. 247 #[derive(Clone, Copy, PartialEq, Eq)] 248 pub struct MycDeliveryStateError { 249 kind: MycDeliveryStateErrorKind, 250 } 251 252 impl MycDeliveryStateError { 253 const fn new(kind: MycDeliveryStateErrorKind) -> Self { 254 Self { kind } 255 } 256 257 /// Returns the stable failure class. 258 #[must_use] 259 pub const fn kind(self) -> MycDeliveryStateErrorKind { 260 self.kind 261 } 262 263 /// Returns the stable machine-readable classification. 264 #[must_use] 265 pub const fn code(self) -> &'static str { 266 self.kind.code() 267 } 268 } 269 270 impl fmt::Display for MycDeliveryStateError { 271 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 272 formatter.write_str(match self.kind { 273 MycDeliveryStateErrorKind::InvalidTime => "delivery time is invalid", 274 MycDeliveryStateErrorKind::InvalidRelayId => "delivery relay identity is invalid", 275 MycDeliveryStateErrorKind::InvalidPolicy => "delivery policy is invalid", 276 MycDeliveryStateErrorKind::InvalidRetryJitter => "delivery retry jitter is invalid", 277 }) 278 } 279 } 280 281 impl fmt::Debug for MycDeliveryStateError { 282 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 283 formatter 284 .debug_struct("MycDeliveryStateError") 285 .field("kind", &self.kind) 286 .finish() 287 } 288 } 289 290 impl Error for MycDeliveryStateError {} 291 292 /// Caller-injected full-jitter delay, later relationship-checked against a job. 293 /// 294 /// The runtime entropy adapter owns generation. State code accepts only this 295 /// bounded value and persists the exact resulting retry schedule. 296 #[derive(Clone, Copy, Debug, PartialEq, Eq, PartialOrd, Ord, Hash)] 297 pub struct MycDeliveryRetryJitter(u64); 298 299 impl MycDeliveryRetryJitter { 300 /// Validates one injected full-jitter delay. 301 pub fn new(value_ms: u64) -> Result<Self, MycDeliveryStateError> { 302 if value_ms > MYC_DELIVERY_RETRY_JITTER_MAX_MS { 303 return Err(MycDeliveryStateError::new( 304 MycDeliveryStateErrorKind::InvalidRetryJitter, 305 )); 306 } 307 Ok(Self(value_ms)) 308 } 309 310 /// Returns the exact injected delay in milliseconds. 311 #[must_use] 312 pub const fn get(self) -> u64 { 313 self.0 314 } 315 } 316 317 macro_rules! redacted_id { 318 ($name:ident, $debug:literal) => { 319 #[derive(Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash)] 320 pub struct $name([u8; 32]); 321 322 impl $name { 323 /// Returns the exact identity bytes. 324 #[must_use] 325 pub const fn as_bytes(&self) -> &[u8; 32] { 326 &self.0 327 } 328 } 329 330 impl fmt::Debug for $name { 331 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 332 formatter.write_str($debug) 333 } 334 } 335 }; 336 } 337 338 redacted_id!(MycDeliveryJobId, "MycDeliveryJobId([redacted])"); 339 redacted_id!(MycDeliveryAttemptId, "MycDeliveryAttemptId([redacted])"); 340 redacted_id!( 341 MycDeliveryArtifactDigest, 342 "MycDeliveryArtifactDigest([redacted])" 343 ); 344 345 impl MycDeliveryJobId { 346 pub(crate) const fn from_persisted(bytes: [u8; 32]) -> Self { 347 Self(bytes) 348 } 349 } 350 351 impl MycDeliveryArtifactDigest { 352 /// Wraps an independently verified exact-artifact SHA-256 identity. 353 #[must_use] 354 pub const fn from_bytes(bytes: [u8; 32]) -> Self { 355 Self(bytes) 356 } 357 } 358 359 /// One-use injected entropy for a new delivery attempt lease. 360 pub struct MycDeliveryAttemptNonce([u8; 32]); 361 362 impl MycDeliveryAttemptNonce { 363 /// Wraps exact entropy supplied by the caller's injected boundary. 364 #[must_use] 365 pub const fn from_injected_entropy(bytes: [u8; 32]) -> Self { 366 Self(bytes) 367 } 368 } 369 370 impl fmt::Debug for MycDeliveryAttemptNonce { 371 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 372 formatter.write_str("MycDeliveryAttemptNonce([redacted])") 373 } 374 } 375 376 /// Positive UTC millisecond evidence representable by SQLite. 377 #[derive(Clone, Copy, Debug, PartialEq, Eq, PartialOrd, Ord, Hash)] 378 pub struct MycDeliveryTimeUnixMs(u64); 379 380 impl MycDeliveryTimeUnixMs { 381 /// Validates one positive UTC millisecond instant. 382 pub fn new(value: u64) -> Result<Self, MycDeliveryStateError> { 383 if value == 0 || i64::try_from(value).is_err() { 384 return Err(MycDeliveryStateError::new( 385 MycDeliveryStateErrorKind::InvalidTime, 386 )); 387 } 388 Ok(Self(value)) 389 } 390 391 /// Returns the validated instant. 392 #[must_use] 393 pub const fn get(self) -> u64 { 394 self.0 395 } 396 397 pub(crate) fn sqlite_value(self) -> i64 { 398 i64::try_from(self.0).expect("validated delivery time fits SQLite") 399 } 400 } 401 402 /// Canonical configured relay identifier used as a durable target identity. 403 #[derive(Clone, PartialEq, Eq, PartialOrd, Ord, Hash)] 404 pub struct MycDeliveryRelayId(Box<str>); 405 406 impl MycDeliveryRelayId { 407 /// Validates the frozen lower-snake relay grammar before allocation. 408 pub fn new(value: &str) -> Result<Self, MycDeliveryStateError> { 409 let bytes = value.as_bytes(); 410 let valid = !bytes.is_empty() 411 && bytes.len() <= MYC_DELIVERY_RELAY_ID_MAX_BYTES 412 && bytes[0].is_ascii_lowercase() 413 && bytes[bytes.len() - 1].is_ascii_alphanumeric() 414 && bytes 415 .iter() 416 .all(|byte| byte.is_ascii_lowercase() || byte.is_ascii_digit() || *byte == b'_') 417 && !bytes.windows(2).any(|window| window == b"__"); 418 if !valid { 419 return Err(MycDeliveryStateError::new( 420 MycDeliveryStateErrorKind::InvalidRelayId, 421 )); 422 } 423 Ok(Self(value.into())) 424 } 425 426 /// Returns the canonical identifier. 427 #[must_use] 428 pub fn as_str(&self) -> &str { 429 &self.0 430 } 431 } 432 433 impl fmt::Debug for MycDeliveryRelayId { 434 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 435 formatter.write_str("MycDeliveryRelayId([redacted])") 436 } 437 } 438 439 /// Closed delivery-policy mode copied from normalized configuration. 440 #[derive(Clone, Copy, Debug, PartialEq, Eq)] 441 pub enum MycDeliveryPolicyMode { 442 AtLeastOneRequired, 443 AllRequired, 444 RequiredQuorum, 445 } 446 447 impl MycDeliveryPolicyMode { 448 /// Returns the exact durable spelling. 449 #[must_use] 450 pub const fn as_str(self) -> &'static str { 451 match self { 452 Self::AtLeastOneRequired => "at_least_one_required", 453 Self::AllRequired => "all_required", 454 Self::RequiredQuorum => "required_quorum", 455 } 456 } 457 458 pub(crate) fn parse(value: &str) -> Option<Self> { 459 match value { 460 "at_least_one_required" => Some(Self::AtLeastOneRequired), 461 "all_required" => Some(Self::AllRequired), 462 "required_quorum" => Some(Self::RequiredQuorum), 463 _ => None, 464 } 465 } 466 } 467 468 /// Immutable configured target identity. 469 #[derive(Clone, PartialEq, Eq)] 470 pub(crate) struct MycDeliveryTargetPolicy { 471 relay_id: MycDeliveryRelayId, 472 required: bool, 473 } 474 475 /// Immutable configured delivery and retry authority. 476 #[derive(Clone, PartialEq, Eq)] 477 pub(crate) struct MycDeliveryPolicies { 478 mode: MycDeliveryPolicyMode, 479 required_acknowledgements: u32, 480 max_attempts: u32, 481 initial_backoff_ms: u64, 482 maximum_backoff_ms: u64, 483 attempt_deadline_ms: u64, 484 outbox_maximum: usize, 485 targets: Box<[MycDeliveryTargetPolicy]>, 486 } 487 488 impl MycDeliveryPolicies { 489 #[allow(clippy::too_many_arguments)] 490 pub(crate) fn new( 491 mode: MycDeliveryPolicyMode, 492 configured_quorum: Option<u32>, 493 max_attempts: u32, 494 initial_backoff_ms: u64, 495 maximum_backoff_ms: u64, 496 attempt_deadline_ms: u64, 497 outbox_maximum: usize, 498 mut targets: Vec<(MycDeliveryRelayId, bool)>, 499 ) -> Result<Self, MycDeliveryStateError> { 500 targets.sort_by(|left, right| left.0.cmp(&right.0)); 501 let required_count = targets.iter().filter(|(_, required)| *required).count(); 502 let required_acknowledgements = match mode { 503 MycDeliveryPolicyMode::AtLeastOneRequired => 1, 504 MycDeliveryPolicyMode::AllRequired => u32::try_from(required_count).unwrap_or(u32::MAX), 505 MycDeliveryPolicyMode::RequiredQuorum => configured_quorum.unwrap_or(0), 506 }; 507 let valid = !targets.is_empty() 508 && targets.len() <= MYC_DELIVERY_TARGET_MAX_COUNT 509 && !targets.windows(2).any(|window| window[0].0 == window[1].0) 510 && required_acknowledgements != 0 511 && usize::try_from(required_acknowledgements) 512 .is_ok_and(|required| required <= targets.len()) 513 && usize::try_from(required_acknowledgements) 514 .is_ok_and(|required| required <= required_count) 515 && (1..=MYC_DELIVERY_ATTEMPT_MAX_COUNT).contains(&max_attempts) 516 && initial_backoff_ms != 0 517 && initial_backoff_ms <= maximum_backoff_ms 518 && maximum_backoff_ms <= 300_000 519 && attempt_deadline_ms != 0 520 && attempt_deadline_ms <= 30_000 521 && (1..=65_536).contains(&outbox_maximum); 522 if !valid { 523 return Err(MycDeliveryStateError::new( 524 MycDeliveryStateErrorKind::InvalidPolicy, 525 )); 526 } 527 Ok(Self { 528 mode, 529 required_acknowledgements, 530 max_attempts, 531 initial_backoff_ms, 532 maximum_backoff_ms, 533 attempt_deadline_ms, 534 outbox_maximum, 535 targets: targets 536 .into_iter() 537 .map(|(relay_id, required)| MycDeliveryTargetPolicy { relay_id, required }) 538 .collect(), 539 }) 540 } 541 542 pub(crate) const fn outbox_maximum(&self) -> usize { 543 self.outbox_maximum 544 } 545 } 546 547 impl fmt::Debug for MycDeliveryPolicies { 548 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 549 formatter 550 .debug_struct("MycDeliveryPolicies") 551 .field("mode", &self.mode) 552 .field("target_count", &self.targets.len()) 553 .finish_non_exhaustive() 554 } 555 } 556 557 /// Closed authority that produced the exact bytes retained by a delivery job. 558 #[derive(Clone, Copy, Debug, PartialEq, Eq)] 559 pub enum MycDeliverySourceKind { 560 SignerResponse, 561 DiscoveryHandler, 562 } 563 564 impl MycDeliverySourceKind { 565 /// Returns the exact durable spelling. 566 #[must_use] 567 pub const fn as_str(self) -> &'static str { 568 match self { 569 Self::SignerResponse => "signer_response", 570 Self::DiscoveryHandler => "discovery_handler", 571 } 572 } 573 574 pub(crate) fn parse(value: &str) -> Option<Self> { 575 match value { 576 "signer_response" => Some(Self::SignerResponse), 577 "discovery_handler" => Some(Self::DiscoveryHandler), 578 _ => None, 579 } 580 } 581 } 582 583 #[derive(Clone, Copy, PartialEq, Eq)] 584 pub(crate) struct MycDeliverySource { 585 kind: MycDeliverySourceKind, 586 id: [u8; 32], 587 } 588 589 impl MycDeliverySource { 590 pub(crate) const fn signer_response(operation_id: MycSignerOperationId) -> Self { 591 Self { 592 kind: MycDeliverySourceKind::SignerResponse, 593 id: *operation_id.as_bytes(), 594 } 595 } 596 597 pub(crate) const fn discovery_handler(generation_id: [u8; 32]) -> Self { 598 Self { 599 kind: MycDeliverySourceKind::DiscoveryHandler, 600 id: generation_id, 601 } 602 } 603 604 pub(crate) const fn kind(self) -> MycDeliverySourceKind { 605 self.kind 606 } 607 608 pub(crate) const fn id(self) -> [u8; 32] { 609 self.id 610 } 611 } 612 613 impl fmt::Debug for MycDeliverySource { 614 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 615 formatter 616 .debug_struct("MycDeliverySource") 617 .field("kind", &self.kind) 618 .field("id", &"[redacted]") 619 .finish() 620 } 621 } 622 623 /// Durable job state. 624 #[derive(Clone, Copy, Debug, PartialEq, Eq)] 625 pub enum MycDeliveryJobStatus { 626 Pending, 627 Active, 628 Delivered, 629 Failed, 630 Unknown, 631 } 632 633 impl MycDeliveryJobStatus { 634 /// Returns the exact durable spelling. 635 #[must_use] 636 pub const fn as_str(self) -> &'static str { 637 match self { 638 Self::Pending => "pending", 639 Self::Active => "active", 640 Self::Delivered => "delivered", 641 Self::Failed => "failed", 642 Self::Unknown => "unknown", 643 } 644 } 645 646 pub(crate) fn parse(value: &str) -> Option<Self> { 647 match value { 648 "pending" => Some(Self::Pending), 649 "active" => Some(Self::Active), 650 "delivered" => Some(Self::Delivered), 651 "failed" => Some(Self::Failed), 652 "unknown" => Some(Self::Unknown), 653 _ => None, 654 } 655 } 656 } 657 658 /// Durable per-target state. 659 #[derive(Clone, Copy, Debug, PartialEq, Eq)] 660 pub enum MycDeliveryTargetStatus { 661 Pending, 662 Leased, 663 Submitted, 664 Delivered, 665 Retryable, 666 Unknown, 667 Exhausted, 668 } 669 670 impl MycDeliveryTargetStatus { 671 /// Returns the exact durable spelling. 672 #[must_use] 673 pub const fn as_str(self) -> &'static str { 674 match self { 675 Self::Pending => "pending", 676 Self::Leased => "leased", 677 Self::Submitted => "submitted", 678 Self::Delivered => "delivered", 679 Self::Retryable => "retryable", 680 Self::Unknown => "unknown", 681 Self::Exhausted => "exhausted", 682 } 683 } 684 685 fn parse(value: &str) -> Option<Self> { 686 match value { 687 "pending" => Some(Self::Pending), 688 "leased" => Some(Self::Leased), 689 "submitted" => Some(Self::Submitted), 690 "delivered" => Some(Self::Delivered), 691 "retryable" => Some(Self::Retryable), 692 "unknown" => Some(Self::Unknown), 693 "exhausted" => Some(Self::Exhausted), 694 _ => None, 695 } 696 } 697 } 698 699 /// Durable attempt state; `Unknown` is distinct from proof of failure. 700 #[derive(Clone, Copy, Debug, PartialEq, Eq)] 701 pub enum MycDeliveryAttemptStatus { 702 Leased, 703 Submitted, 704 Delivered, 705 Failed, 706 Unknown, 707 } 708 709 impl MycDeliveryAttemptStatus { 710 /// Returns the exact durable spelling. 711 #[must_use] 712 pub const fn as_str(self) -> &'static str { 713 match self { 714 Self::Leased => "leased", 715 Self::Submitted => "submitted", 716 Self::Delivered => "delivered", 717 Self::Failed => "failed", 718 Self::Unknown => "unknown", 719 } 720 } 721 722 fn parse(value: &str) -> Option<Self> { 723 match value { 724 "leased" => Some(Self::Leased), 725 "submitted" => Some(Self::Submitted), 726 "delivered" => Some(Self::Delivered), 727 "failed" => Some(Self::Failed), 728 "unknown" => Some(Self::Unknown), 729 _ => None, 730 } 731 } 732 } 733 734 /// Closed evidence for a terminal attempt observation. 735 #[derive(Clone, Copy, Debug, PartialEq, Eq)] 736 pub enum MycDeliveryAttemptOutcome { 737 Delivered, 738 RelayRejected, 739 TransportFailed, 740 UnknownAcknowledgement, 741 } 742 743 impl MycDeliveryAttemptOutcome { 744 const fn status(self) -> MycDeliveryAttemptStatus { 745 match self { 746 Self::Delivered => MycDeliveryAttemptStatus::Delivered, 747 Self::RelayRejected | Self::TransportFailed => MycDeliveryAttemptStatus::Failed, 748 Self::UnknownAcknowledgement => MycDeliveryAttemptStatus::Unknown, 749 } 750 } 751 752 const fn reason(self) -> &'static str { 753 match self { 754 Self::Delivered => "accepted", 755 Self::RelayRejected => "relay_rejected", 756 Self::TransportFailed => "transport_failed", 757 Self::UnknownAcknowledgement => "acknowledgement_lost", 758 } 759 } 760 } 761 762 /// Immutable summary of one retained delivery job. 763 #[derive(Clone, PartialEq, Eq)] 764 pub struct MycDeliveryJobRecord { 765 id: MycDeliveryJobId, 766 source: MycDeliverySource, 767 artifact_digest: MycDeliveryArtifactDigest, 768 policy_mode: MycDeliveryPolicyMode, 769 required_acknowledgements: u32, 770 max_attempts: u32, 771 initial_backoff_ms: u64, 772 maximum_backoff_ms: u64, 773 attempt_deadline_ms: u64, 774 status: MycDeliveryJobStatus, 775 created_at: MycDeliveryTimeUnixMs, 776 updated_at: MycDeliveryTimeUnixMs, 777 finalized_at: Option<MycDeliveryTimeUnixMs>, 778 targets: Box<[MycDeliveryTargetRecord]>, 779 } 780 781 impl MycDeliveryJobRecord { 782 pub(crate) const fn source_id(&self) -> &[u8; 32] { 783 &self.source.id 784 } 785 786 #[must_use] 787 pub const fn id(&self) -> MycDeliveryJobId { 788 self.id 789 } 790 #[must_use] 791 pub const fn operation_id(&self) -> Option<MycSignerOperationId> { 792 match self.source.kind { 793 MycDeliverySourceKind::SignerResponse => { 794 Some(MycSignerOperationId::from_persisted(self.source.id)) 795 } 796 MycDeliverySourceKind::DiscoveryHandler => None, 797 } 798 } 799 #[must_use] 800 pub const fn source_kind(&self) -> MycDeliverySourceKind { 801 self.source.kind 802 } 803 #[must_use] 804 pub const fn artifact_digest(&self) -> MycDeliveryArtifactDigest { 805 self.artifact_digest 806 } 807 #[must_use] 808 pub const fn policy_mode(&self) -> MycDeliveryPolicyMode { 809 self.policy_mode 810 } 811 #[must_use] 812 pub const fn required_acknowledgements(&self) -> u32 { 813 self.required_acknowledgements 814 } 815 #[must_use] 816 pub const fn max_attempts(&self) -> u32 { 817 self.max_attempts 818 } 819 #[must_use] 820 pub const fn initial_backoff_ms(&self) -> u64 { 821 self.initial_backoff_ms 822 } 823 #[must_use] 824 pub const fn maximum_backoff_ms(&self) -> u64 { 825 self.maximum_backoff_ms 826 } 827 #[must_use] 828 pub const fn attempt_deadline_ms(&self) -> u64 { 829 self.attempt_deadline_ms 830 } 831 #[must_use] 832 pub const fn status(&self) -> MycDeliveryJobStatus { 833 self.status 834 } 835 #[must_use] 836 pub const fn targets(&self) -> &[MycDeliveryTargetRecord] { 837 &self.targets 838 } 839 #[must_use] 840 pub const fn created_at(&self) -> MycDeliveryTimeUnixMs { 841 self.created_at 842 } 843 #[must_use] 844 pub const fn updated_at(&self) -> MycDeliveryTimeUnixMs { 845 self.updated_at 846 } 847 #[must_use] 848 pub const fn finalized_at(&self) -> Option<MycDeliveryTimeUnixMs> { 849 self.finalized_at 850 } 851 } 852 853 impl fmt::Debug for MycDeliveryJobRecord { 854 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 855 formatter 856 .debug_struct("MycDeliveryJobRecord") 857 .field("source_kind", &self.source.kind) 858 .field("status", &self.status) 859 .field("target_count", &self.targets.len()) 860 .finish() 861 } 862 } 863 864 /// Immutable summary of one configured target and its current state. 865 #[derive(Clone, PartialEq, Eq)] 866 pub struct MycDeliveryTargetRecord { 867 index: u32, 868 relay_id: MycDeliveryRelayId, 869 required: bool, 870 attempt_count: u32, 871 status: MycDeliveryTargetStatus, 872 active_attempt_id: Option<MycDeliveryAttemptId>, 873 next_attempt_at: Option<MycDeliveryTimeUnixMs>, 874 updated_at: MycDeliveryTimeUnixMs, 875 } 876 877 impl MycDeliveryTargetRecord { 878 #[must_use] 879 pub const fn index(&self) -> u32 { 880 self.index 881 } 882 #[must_use] 883 pub const fn relay_id(&self) -> &MycDeliveryRelayId { 884 &self.relay_id 885 } 886 #[must_use] 887 pub const fn required(&self) -> bool { 888 self.required 889 } 890 #[must_use] 891 pub const fn attempt_count(&self) -> u32 { 892 self.attempt_count 893 } 894 #[must_use] 895 pub const fn status(&self) -> MycDeliveryTargetStatus { 896 self.status 897 } 898 #[must_use] 899 pub const fn active_attempt_id(&self) -> Option<MycDeliveryAttemptId> { 900 self.active_attempt_id 901 } 902 #[must_use] 903 pub const fn next_attempt_at(&self) -> Option<MycDeliveryTimeUnixMs> { 904 self.next_attempt_at 905 } 906 #[must_use] 907 pub const fn updated_at(&self) -> MycDeliveryTimeUnixMs { 908 self.updated_at 909 } 910 } 911 912 impl fmt::Debug for MycDeliveryTargetRecord { 913 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 914 formatter 915 .debug_struct("MycDeliveryTargetRecord") 916 .field("required", &self.required) 917 .field("attempt_count", &self.attempt_count) 918 .field("status", &self.status) 919 .finish() 920 } 921 } 922 923 /// Immutable identity plus append-only state for one bounded delivery attempt. 924 #[derive(Clone, PartialEq, Eq)] 925 pub struct MycDeliveryAttemptRecord { 926 id: MycDeliveryAttemptId, 927 number: u32, 928 status: MycDeliveryAttemptStatus, 929 leased_at: MycDeliveryTimeUnixMs, 930 lease_expires_at: MycDeliveryTimeUnixMs, 931 submitted_at: Option<MycDeliveryTimeUnixMs>, 932 resolved_at: Option<MycDeliveryTimeUnixMs>, 933 reason: Option<&'static str>, 934 } 935 936 impl MycDeliveryAttemptRecord { 937 #[must_use] 938 pub const fn id(&self) -> MycDeliveryAttemptId { 939 self.id 940 } 941 #[must_use] 942 pub const fn number(&self) -> u32 { 943 self.number 944 } 945 #[must_use] 946 pub const fn status(&self) -> MycDeliveryAttemptStatus { 947 self.status 948 } 949 #[must_use] 950 pub const fn leased_at(&self) -> MycDeliveryTimeUnixMs { 951 self.leased_at 952 } 953 #[must_use] 954 pub const fn lease_expires_at(&self) -> MycDeliveryTimeUnixMs { 955 self.lease_expires_at 956 } 957 #[must_use] 958 pub const fn submitted_at(&self) -> Option<MycDeliveryTimeUnixMs> { 959 self.submitted_at 960 } 961 #[must_use] 962 pub const fn resolved_at(&self) -> Option<MycDeliveryTimeUnixMs> { 963 self.resolved_at 964 } 965 #[must_use] 966 pub const fn reason(&self) -> Option<&'static str> { 967 self.reason 968 } 969 } 970 971 impl fmt::Debug for MycDeliveryAttemptRecord { 972 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 973 formatter 974 .debug_struct("MycDeliveryAttemptRecord") 975 .field("number", &self.number) 976 .field("status", &self.status) 977 .finish() 978 } 979 } 980 981 /// New or idempotently replayed delivery-job creation. 982 #[derive(Clone, PartialEq, Eq)] 983 pub enum MycDeliveryJobAdmission { 984 Created(MycDeliveryJobRecord), 985 ExactReplay(MycDeliveryJobRecord), 986 } 987 988 impl MycDeliveryJobAdmission { 989 #[must_use] 990 pub const fn record(&self) -> &MycDeliveryJobRecord { 991 match self { 992 Self::Created(record) | Self::ExactReplay(record) => record, 993 } 994 } 995 } 996 997 impl fmt::Debug for MycDeliveryJobAdmission { 998 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 999 formatter.write_str(match self { 1000 Self::Created(_) => "MycDeliveryJobAdmission::Created([redacted])", 1001 Self::ExactReplay(_) => "MycDeliveryJobAdmission::ExactReplay([redacted])", 1002 }) 1003 } 1004 } 1005 1006 /// Result of attempting to claim one target. 1007 #[derive(Clone, PartialEq, Eq)] 1008 pub enum MycDeliveryClaim { 1009 Claimed(MycDeliveryAttemptRecord), 1010 ExactReplay(MycDeliveryAttemptRecord), 1011 NotReady, 1012 Terminal, 1013 } 1014 1015 impl fmt::Debug for MycDeliveryClaim { 1016 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 1017 formatter.write_str(match self { 1018 Self::Claimed(_) => "MycDeliveryClaim::Claimed([redacted])", 1019 Self::ExactReplay(_) => "MycDeliveryClaim::ExactReplay([redacted])", 1020 Self::NotReady => "MycDeliveryClaim::NotReady", 1021 Self::Terminal => "MycDeliveryClaim::Terminal", 1022 }) 1023 } 1024 } 1025 1026 impl MycStateRepository<'_> { 1027 pub(crate) async fn read_runtime_outbox_status( 1028 &self, 1029 ) -> Result<crate::MycOutboxStatusV1, MycStateRepositoryError> { 1030 let expected = PersistedMetadata::from(self.expected()); 1031 self.host() 1032 .transaction(move |transaction| { 1033 Box::pin(async move { 1034 verify_metadata(transaction, &expected).await?; 1035 let rows = sqlx::query(READ_RUNTIME_OUTBOX_STATUS_SQL) 1036 .fetch_all(&mut *transaction) 1037 .await 1038 .map_err(|_| DeliveryOperationError::Storage)?; 1039 let [row] = rows.as_slice() else { 1040 return Err(DeliveryOperationError::Binding); 1041 }; 1042 let count = |column| { 1043 row.try_get::<i64, _>(column) 1044 .ok() 1045 .and_then(|value| u64::try_from(value).ok()) 1046 .ok_or(DeliveryOperationError::Binding) 1047 }; 1048 let oldest = match row 1049 .try_get::<&str, _>("oldest_pending_type") 1050 .map_err(|_| DeliveryOperationError::Binding)? 1051 { 1052 "null" => None, 1053 "integer" => { 1054 let milliseconds = count("oldest_pending_at_unix_ms")?; 1055 Some( 1056 crate::MycStatusUnixSeconds::new(milliseconds / 1_000) 1057 .map_err(|_| DeliveryOperationError::Binding)?, 1058 ) 1059 } 1060 _ => return Err(DeliveryOperationError::Binding), 1061 }; 1062 Ok(crate::MycOutboxStatusV1::new( 1063 count("pending_count")?, 1064 count("unknown_count")?, 1065 oldest, 1066 )) 1067 }) 1068 }) 1069 .await 1070 .map_err(map_transaction_error) 1071 } 1072 1073 /// Returns the first exact target eligible for bounded delivery work. 1074 /// 1075 /// Selection is deterministic and performs no claim or network I/O. The 1076 /// subsequent claim transaction remains the sole lease authority. 1077 pub(crate) async fn next_ready_delivery_target( 1078 &self, 1079 observed_at: MycDeliveryTimeUnixMs, 1080 ) -> Result<Option<(MycDeliveryJobId, MycDeliveryRelayId)>, MycStateRepositoryError> { 1081 let expected = PersistedMetadata::from(self.expected()); 1082 self.host() 1083 .transaction(move |transaction| { 1084 Box::pin(async move { 1085 verify_metadata(transaction, &expected).await?; 1086 let rows = sqlx::query(READ_NEXT_READY_TARGET_SQL) 1087 .bind(observed_at.sqlite_value()) 1088 .fetch_all(&mut *transaction) 1089 .await 1090 .map_err(|_| DeliveryOperationError::Storage)?; 1091 match rows.as_slice() { 1092 [] => Ok(None), 1093 [row] => { 1094 let job_id = MycDeliveryJobId::from_persisted(blob32(row, "job_id")?); 1095 let relay_id = MycDeliveryRelayId::new(text(row, "relay_id")?) 1096 .map_err(|_| DeliveryOperationError::Binding)?; 1097 Ok(Some((job_id, relay_id))) 1098 } 1099 _ => Err(DeliveryOperationError::Binding), 1100 } 1101 }) 1102 }) 1103 .await 1104 .map_err(map_transaction_error) 1105 } 1106 1107 /// Claims one eligible target under a bounded expiring attempt lease. 1108 pub async fn claim_delivery_target( 1109 &self, 1110 job_id: MycDeliveryJobId, 1111 relay_id: &MycDeliveryRelayId, 1112 nonce: MycDeliveryAttemptNonce, 1113 claimed_at: MycDeliveryTimeUnixMs, 1114 ) -> Result<MycDeliveryClaim, MycStateRepositoryError> { 1115 let relay_id = relay_id.clone(); 1116 let expected = PersistedMetadata::from(self.expected()); 1117 self.host() 1118 .transaction(move |transaction| { 1119 Box::pin(async move { 1120 verify_metadata(transaction, &expected).await?; 1121 claim_target(transaction, job_id, &relay_id, &nonce, claimed_at).await 1122 }) 1123 }) 1124 .await 1125 .map_err(map_transaction_error) 1126 } 1127 1128 /// Records that the exact leased attempt reached the external submission boundary. 1129 pub async fn mark_delivery_attempt_submitted( 1130 &self, 1131 job_id: MycDeliveryJobId, 1132 relay_id: &MycDeliveryRelayId, 1133 attempt_id: MycDeliveryAttemptId, 1134 submitted_at: MycDeliveryTimeUnixMs, 1135 ) -> Result<MycDeliveryAttemptRecord, MycStateRepositoryError> { 1136 let relay_id = relay_id.clone(); 1137 let expected = PersistedMetadata::from(self.expected()); 1138 self.host() 1139 .transaction(move |transaction| { 1140 Box::pin(async move { 1141 verify_metadata(transaction, &expected).await?; 1142 mark_submitted(transaction, job_id, &relay_id, attempt_id, submitted_at).await 1143 }) 1144 }) 1145 .await 1146 .map_err(map_transaction_error) 1147 } 1148 1149 /// Persists delivered, proven-failed, or unknown acknowledgement evidence. 1150 pub async fn record_delivery_attempt_outcome( 1151 &self, 1152 job_id: MycDeliveryJobId, 1153 relay_id: &MycDeliveryRelayId, 1154 attempt_id: MycDeliveryAttemptId, 1155 outcome: MycDeliveryAttemptOutcome, 1156 retry_jitter: MycDeliveryRetryJitter, 1157 observed_at: MycDeliveryTimeUnixMs, 1158 ) -> Result<MycDeliveryJobRecord, MycStateRepositoryError> { 1159 let relay_id = relay_id.clone(); 1160 let expected = PersistedMetadata::from(self.expected()); 1161 self.host() 1162 .transaction(move |transaction| { 1163 Box::pin(async move { 1164 verify_metadata(transaction, &expected).await?; 1165 record_outcome( 1166 transaction, 1167 job_id, 1168 &relay_id, 1169 attempt_id, 1170 outcome, 1171 retry_jitter, 1172 observed_at, 1173 ) 1174 .await 1175 }) 1176 }) 1177 .await 1178 .map_err(map_transaction_error) 1179 } 1180 1181 /// Converts an expired pre-submit lease to failure or a submitted lease to unknown. 1182 pub async fn recover_expired_delivery_lease( 1183 &self, 1184 job_id: MycDeliveryJobId, 1185 relay_id: &MycDeliveryRelayId, 1186 attempt_id: MycDeliveryAttemptId, 1187 retry_jitter: MycDeliveryRetryJitter, 1188 observed_at: MycDeliveryTimeUnixMs, 1189 ) -> Result<MycDeliveryJobRecord, MycStateRepositoryError> { 1190 let relay_id = relay_id.clone(); 1191 let expected = PersistedMetadata::from(self.expected()); 1192 self.host() 1193 .transaction(move |transaction| { 1194 Box::pin(async move { 1195 verify_metadata(transaction, &expected).await?; 1196 recover_expired( 1197 transaction, 1198 job_id, 1199 &relay_id, 1200 attempt_id, 1201 retry_jitter, 1202 observed_at, 1203 ) 1204 .await 1205 }) 1206 }) 1207 .await 1208 .map_err(map_transaction_error) 1209 } 1210 1211 /// Reads one bounded job snapshot and its immutable target set. 1212 pub async fn read_delivery_job( 1213 &self, 1214 job_id: MycDeliveryJobId, 1215 ) -> Result<Option<MycDeliveryJobRecord>, MycStateRepositoryError> { 1216 let expected = PersistedMetadata::from(self.expected()); 1217 self.host() 1218 .transaction(move |transaction| { 1219 Box::pin(async move { 1220 verify_metadata(transaction, &expected).await?; 1221 read_job(transaction, job_id).await 1222 }) 1223 }) 1224 .await 1225 .map_err(map_transaction_error) 1226 } 1227 1228 /// Reads at most the configured 32 attempts for one retained target. 1229 pub async fn read_delivery_attempts( 1230 &self, 1231 job_id: MycDeliveryJobId, 1232 relay_id: &MycDeliveryRelayId, 1233 ) -> Result<Box<[MycDeliveryAttemptRecord]>, MycStateRepositoryError> { 1234 let relay_id = relay_id.clone(); 1235 let expected = PersistedMetadata::from(self.expected()); 1236 self.host() 1237 .transaction(move |transaction| { 1238 Box::pin(async move { 1239 verify_metadata(transaction, &expected).await?; 1240 let job = read_job(transaction, job_id) 1241 .await? 1242 .ok_or(DeliveryOperationError::Binding)?; 1243 let target = target_by_relay(&job, &relay_id)?; 1244 read_attempts(transaction, job_id, target.index).await 1245 }) 1246 }) 1247 .await 1248 .map_err(map_transaction_error) 1249 } 1250 } 1251 1252 #[derive(Clone, Copy, Debug, PartialEq, Eq)] 1253 pub(crate) enum DeliveryOperationError { 1254 Binding, 1255 Storage, 1256 } 1257 1258 async fn verify_metadata( 1259 transaction: &mut ServiceSqliteTransaction<'_>, 1260 expected: &PersistedMetadata, 1261 ) -> Result<(), DeliveryOperationError> { 1262 require_expected_metadata(transaction, expected) 1263 .await 1264 .map_err(|error| match error { 1265 RepositoryOperationError::Binding => DeliveryOperationError::Binding, 1266 RepositoryOperationError::Storage => DeliveryOperationError::Storage, 1267 }) 1268 } 1269 1270 pub(crate) async fn create_job( 1271 transaction: &mut ServiceSqliteTransaction<'_>, 1272 source: MycDeliverySource, 1273 artifact_digest: MycDeliveryArtifactDigest, 1274 created_at: MycDeliveryTimeUnixMs, 1275 policy: &MycDeliveryPolicies, 1276 ) -> Result<MycDeliveryJobAdmission, DeliveryOperationError> { 1277 if let Some(existing) = read_job_by_source(transaction, source).await? { 1278 return exact_job(&existing, source, artifact_digest, created_at, policy) 1279 .then_some(MycDeliveryJobAdmission::ExactReplay(existing)) 1280 .ok_or(DeliveryOperationError::Binding); 1281 } 1282 let active_jobs = sqlx::query(READ_ACTIVE_JOB_COUNT_SQL) 1283 .fetch_one(&mut *transaction) 1284 .await 1285 .map_err(|_| DeliveryOperationError::Storage)? 1286 .try_get::<i64, _>("row_count") 1287 .map_err(|_| DeliveryOperationError::Binding)?; 1288 if usize::try_from(active_jobs) 1289 .ok() 1290 .is_none_or(|count| count >= policy.outbox_maximum) 1291 { 1292 return Err(DeliveryOperationError::Binding); 1293 } 1294 let job_id = derive_job_id(source, artifact_digest); 1295 let result = sqlx::query(INSERT_JOB_SQL) 1296 .bind(job_id.as_bytes().as_slice()) 1297 .bind(source.kind().as_str()) 1298 .bind(source.id().as_slice()) 1299 .bind(artifact_digest.as_bytes().as_slice()) 1300 .bind(policy.mode.as_str()) 1301 .bind(i64::from(policy.required_acknowledgements)) 1302 .bind(i64::from(policy.max_attempts)) 1303 .bind( 1304 i64::try_from(policy.initial_backoff_ms) 1305 .map_err(|_| DeliveryOperationError::Binding)?, 1306 ) 1307 .bind( 1308 i64::try_from(policy.maximum_backoff_ms) 1309 .map_err(|_| DeliveryOperationError::Binding)?, 1310 ) 1311 .bind( 1312 i64::try_from(policy.attempt_deadline_ms) 1313 .map_err(|_| DeliveryOperationError::Binding)?, 1314 ) 1315 .bind(created_at.sqlite_value()) 1316 .bind(created_at.sqlite_value()) 1317 .execute(&mut *transaction) 1318 .await 1319 .map_err(|_| DeliveryOperationError::Storage)?; 1320 require_one(result.rows_affected())?; 1321 for (index, target) in policy.targets.iter().enumerate() { 1322 let index = u32::try_from(index).map_err(|_| DeliveryOperationError::Binding)?; 1323 let result = sqlx::query(INSERT_TARGET_SQL) 1324 .bind(job_id.as_bytes().as_slice()) 1325 .bind(i64::from(index)) 1326 .bind(target.relay_id.as_str()) 1327 .bind(target.required) 1328 .bind(created_at.sqlite_value()) 1329 .execute(&mut *transaction) 1330 .await 1331 .map_err(|_| DeliveryOperationError::Storage)?; 1332 require_one(result.rows_affected())?; 1333 } 1334 let record = read_job(transaction, job_id) 1335 .await? 1336 .ok_or(DeliveryOperationError::Binding)?; 1337 Ok(MycDeliveryJobAdmission::Created(record)) 1338 } 1339 1340 async fn claim_target( 1341 transaction: &mut ServiceSqliteTransaction<'_>, 1342 job_id: MycDeliveryJobId, 1343 relay_id: &MycDeliveryRelayId, 1344 nonce: &MycDeliveryAttemptNonce, 1345 claimed_at: MycDeliveryTimeUnixMs, 1346 ) -> Result<MycDeliveryClaim, DeliveryOperationError> { 1347 let job = read_job(transaction, job_id) 1348 .await? 1349 .ok_or(DeliveryOperationError::Binding)?; 1350 if matches!( 1351 job.status, 1352 MycDeliveryJobStatus::Delivered 1353 | MycDeliveryJobStatus::Failed 1354 | MycDeliveryJobStatus::Unknown 1355 ) { 1356 return Ok(MycDeliveryClaim::Terminal); 1357 } 1358 let target = target_by_relay(&job, relay_id)?.clone(); 1359 if let Some(existing) = read_attempt_by_nonce(transaction, job_id, target.index, nonce).await? { 1360 return Ok(MycDeliveryClaim::ExactReplay(existing)); 1361 } 1362 if target.active_attempt_id.is_some() 1363 || target.next_attempt_at.is_some_and(|next| next > claimed_at) 1364 { 1365 return Ok(MycDeliveryClaim::NotReady); 1366 } 1367 if matches!( 1368 target.status, 1369 MycDeliveryTargetStatus::Delivered | MycDeliveryTargetStatus::Exhausted 1370 ) || target.attempt_count >= job.max_attempts 1371 { 1372 return Ok(MycDeliveryClaim::Terminal); 1373 } 1374 let attempt_number = target.attempt_count + 1; 1375 let lease_expires_value = claimed_at 1376 .get() 1377 .checked_add(job.attempt_deadline_ms) 1378 .ok_or(DeliveryOperationError::Binding)?; 1379 let lease_expires = MycDeliveryTimeUnixMs::new(lease_expires_value) 1380 .map_err(|_| DeliveryOperationError::Binding)?; 1381 let attempt_id = derive_attempt_id(job_id, target.index, attempt_number, nonce); 1382 let result = sqlx::query(INSERT_ATTEMPT_SQL) 1383 .bind(attempt_id.as_bytes().as_slice()) 1384 .bind(job_id.as_bytes().as_slice()) 1385 .bind(i64::from(target.index)) 1386 .bind(i64::from(attempt_number)) 1387 .bind(nonce.0.as_slice()) 1388 .bind(claimed_at.sqlite_value()) 1389 .bind(lease_expires.sqlite_value()) 1390 .execute(&mut *transaction) 1391 .await 1392 .map_err(|_| DeliveryOperationError::Storage)?; 1393 require_one(result.rows_affected())?; 1394 let result = sqlx::query(CLAIM_TARGET_SQL) 1395 .bind(i64::from(attempt_number)) 1396 .bind(attempt_id.as_bytes().as_slice()) 1397 .bind(claimed_at.sqlite_value()) 1398 .bind(job_id.as_bytes().as_slice()) 1399 .bind(i64::from(target.index)) 1400 .bind(i64::from(target.attempt_count)) 1401 .execute(&mut *transaction) 1402 .await 1403 .map_err(|_| DeliveryOperationError::Storage)?; 1404 require_one(result.rows_affected())?; 1405 let result = sqlx::query(MARK_JOB_ACTIVE_SQL) 1406 .bind(claimed_at.sqlite_value()) 1407 .bind(job_id.as_bytes().as_slice()) 1408 .execute(&mut *transaction) 1409 .await 1410 .map_err(|_| DeliveryOperationError::Storage)?; 1411 if result.rows_affected() > 1 { 1412 return Err(DeliveryOperationError::Storage); 1413 } 1414 let attempt = read_attempt(transaction, job_id, target.index, attempt_id) 1415 .await? 1416 .ok_or(DeliveryOperationError::Binding)?; 1417 Ok(MycDeliveryClaim::Claimed(attempt)) 1418 } 1419 1420 async fn mark_submitted( 1421 transaction: &mut ServiceSqliteTransaction<'_>, 1422 job_id: MycDeliveryJobId, 1423 relay_id: &MycDeliveryRelayId, 1424 attempt_id: MycDeliveryAttemptId, 1425 submitted_at: MycDeliveryTimeUnixMs, 1426 ) -> Result<MycDeliveryAttemptRecord, DeliveryOperationError> { 1427 let job = read_job(transaction, job_id) 1428 .await? 1429 .ok_or(DeliveryOperationError::Binding)?; 1430 let target = target_by_relay(&job, relay_id)?; 1431 let attempt = read_attempt(transaction, job_id, target.index, attempt_id) 1432 .await? 1433 .ok_or(DeliveryOperationError::Binding)?; 1434 if attempt.status == MycDeliveryAttemptStatus::Submitted { 1435 return (attempt.submitted_at == Some(submitted_at)) 1436 .then_some(attempt) 1437 .ok_or(DeliveryOperationError::Binding); 1438 } 1439 if attempt.status != MycDeliveryAttemptStatus::Leased 1440 || submitted_at < attempt.leased_at 1441 || submitted_at > attempt.lease_expires_at 1442 || target.active_attempt_id != Some(attempt_id) 1443 || target.status != MycDeliveryTargetStatus::Leased 1444 { 1445 return Err(DeliveryOperationError::Binding); 1446 } 1447 let result = sqlx::query(MARK_ATTEMPT_SUBMITTED_SQL) 1448 .bind(submitted_at.sqlite_value()) 1449 .bind(attempt_id.as_bytes().as_slice()) 1450 .bind(job_id.as_bytes().as_slice()) 1451 .bind(i64::from(target.index)) 1452 .bind(submitted_at.sqlite_value()) 1453 .execute(&mut *transaction) 1454 .await 1455 .map_err(|_| DeliveryOperationError::Storage)?; 1456 require_one(result.rows_affected())?; 1457 let result = sqlx::query(MARK_TARGET_SUBMITTED_SQL) 1458 .bind(submitted_at.sqlite_value()) 1459 .bind(job_id.as_bytes().as_slice()) 1460 .bind(i64::from(target.index)) 1461 .bind(attempt_id.as_bytes().as_slice()) 1462 .execute(&mut *transaction) 1463 .await 1464 .map_err(|_| DeliveryOperationError::Storage)?; 1465 require_one(result.rows_affected())?; 1466 read_attempt(transaction, job_id, target.index, attempt_id) 1467 .await? 1468 .ok_or(DeliveryOperationError::Binding) 1469 } 1470 1471 async fn record_outcome( 1472 transaction: &mut ServiceSqliteTransaction<'_>, 1473 job_id: MycDeliveryJobId, 1474 relay_id: &MycDeliveryRelayId, 1475 attempt_id: MycDeliveryAttemptId, 1476 outcome: MycDeliveryAttemptOutcome, 1477 retry_jitter: MycDeliveryRetryJitter, 1478 observed_at: MycDeliveryTimeUnixMs, 1479 ) -> Result<MycDeliveryJobRecord, DeliveryOperationError> { 1480 let job = read_job(transaction, job_id) 1481 .await? 1482 .ok_or(DeliveryOperationError::Binding)?; 1483 let target = target_by_relay(&job, relay_id)?.clone(); 1484 let attempt = read_attempt(transaction, job_id, target.index, attempt_id) 1485 .await? 1486 .ok_or(DeliveryOperationError::Binding)?; 1487 if matches!( 1488 attempt.status, 1489 MycDeliveryAttemptStatus::Delivered 1490 | MycDeliveryAttemptStatus::Failed 1491 | MycDeliveryAttemptStatus::Unknown 1492 ) { 1493 return (attempt.status == outcome.status() 1494 && attempt.reason == Some(outcome.reason()) 1495 && attempt.resolved_at == Some(observed_at)) 1496 .then_some(job) 1497 .ok_or(DeliveryOperationError::Binding); 1498 } 1499 let required_prior = match outcome { 1500 MycDeliveryAttemptOutcome::Delivered 1501 | MycDeliveryAttemptOutcome::UnknownAcknowledgement => MycDeliveryAttemptStatus::Submitted, 1502 MycDeliveryAttemptOutcome::RelayRejected => MycDeliveryAttemptStatus::Submitted, 1503 MycDeliveryAttemptOutcome::TransportFailed => attempt.status, 1504 }; 1505 if attempt.status != required_prior 1506 || !matches!( 1507 required_prior, 1508 MycDeliveryAttemptStatus::Leased | MycDeliveryAttemptStatus::Submitted 1509 ) 1510 || target.active_attempt_id != Some(attempt_id) 1511 || observed_at < attempt.leased_at 1512 || attempt 1513 .submitted_at 1514 .is_some_and(|submitted_at| observed_at < submitted_at) 1515 || observed_at > attempt.lease_expires_at 1516 { 1517 return Err(DeliveryOperationError::Binding); 1518 } 1519 resolve_attempt_and_target( 1520 transaction, 1521 &job, 1522 &target, 1523 &attempt, 1524 required_prior, 1525 outcome.status(), 1526 outcome.reason(), 1527 retry_jitter, 1528 observed_at, 1529 ) 1530 .await?; 1531 finalize_job_if_terminal(transaction, job_id, observed_at).await 1532 } 1533 1534 pub(crate) async fn recover_expired( 1535 transaction: &mut ServiceSqliteTransaction<'_>, 1536 job_id: MycDeliveryJobId, 1537 relay_id: &MycDeliveryRelayId, 1538 attempt_id: MycDeliveryAttemptId, 1539 retry_jitter: MycDeliveryRetryJitter, 1540 observed_at: MycDeliveryTimeUnixMs, 1541 ) -> Result<MycDeliveryJobRecord, DeliveryOperationError> { 1542 let job = read_job(transaction, job_id) 1543 .await? 1544 .ok_or(DeliveryOperationError::Binding)?; 1545 let target = target_by_relay(&job, relay_id)?.clone(); 1546 let attempt = read_attempt(transaction, job_id, target.index, attempt_id) 1547 .await? 1548 .ok_or(DeliveryOperationError::Binding)?; 1549 if matches!( 1550 attempt.status, 1551 MycDeliveryAttemptStatus::Delivered 1552 | MycDeliveryAttemptStatus::Failed 1553 | MycDeliveryAttemptStatus::Unknown 1554 ) { 1555 return Ok(job); 1556 } 1557 if observed_at <= attempt.lease_expires_at || target.active_attempt_id != Some(attempt_id) { 1558 return Err(DeliveryOperationError::Binding); 1559 } 1560 let (terminal, reason) = match attempt.status { 1561 MycDeliveryAttemptStatus::Leased => ( 1562 MycDeliveryAttemptStatus::Failed, 1563 "lease_expired_before_submit", 1564 ), 1565 MycDeliveryAttemptStatus::Submitted => { 1566 (MycDeliveryAttemptStatus::Unknown, "acknowledgement_lost") 1567 } 1568 MycDeliveryAttemptStatus::Delivered 1569 | MycDeliveryAttemptStatus::Failed 1570 | MycDeliveryAttemptStatus::Unknown => unreachable!(), 1571 }; 1572 resolve_attempt_and_target( 1573 transaction, 1574 &job, 1575 &target, 1576 &attempt, 1577 attempt.status, 1578 terminal, 1579 reason, 1580 retry_jitter, 1581 observed_at, 1582 ) 1583 .await?; 1584 finalize_job_if_terminal(transaction, job_id, observed_at).await 1585 } 1586 1587 #[allow(clippy::too_many_arguments)] 1588 async fn resolve_attempt_and_target( 1589 transaction: &mut ServiceSqliteTransaction<'_>, 1590 job: &MycDeliveryJobRecord, 1591 target: &MycDeliveryTargetRecord, 1592 attempt: &MycDeliveryAttemptRecord, 1593 prior: MycDeliveryAttemptStatus, 1594 terminal: MycDeliveryAttemptStatus, 1595 reason: &'static str, 1596 retry_jitter: MycDeliveryRetryJitter, 1597 observed_at: MycDeliveryTimeUnixMs, 1598 ) -> Result<(), DeliveryOperationError> { 1599 let result = sqlx::query(RESOLVE_ATTEMPT_SQL) 1600 .bind(terminal.as_str()) 1601 .bind(observed_at.sqlite_value()) 1602 .bind(reason) 1603 .bind(attempt.id.as_bytes().as_slice()) 1604 .bind(job.id.as_bytes().as_slice()) 1605 .bind(i64::from(target.index)) 1606 .bind(prior.as_str()) 1607 .execute(&mut *transaction) 1608 .await 1609 .map_err(|_| DeliveryOperationError::Storage)?; 1610 require_one(result.rows_affected())?; 1611 let attempts_remaining = target.attempt_count < job.max_attempts; 1612 let schedules_retry = attempts_remaining 1613 && matches!( 1614 terminal, 1615 MycDeliveryAttemptStatus::Failed | MycDeliveryAttemptStatus::Unknown 1616 ); 1617 if !schedules_retry && retry_jitter.get() != 0 { 1618 return Err(DeliveryOperationError::Binding); 1619 } 1620 let (target_status, next_attempt) = match terminal { 1621 MycDeliveryAttemptStatus::Delivered => (MycDeliveryTargetStatus::Delivered, None), 1622 MycDeliveryAttemptStatus::Failed if attempts_remaining => ( 1623 MycDeliveryTargetStatus::Retryable, 1624 Some(next_attempt_time( 1625 job, 1626 attempt.number, 1627 retry_jitter, 1628 observed_at, 1629 )?), 1630 ), 1631 MycDeliveryAttemptStatus::Unknown => ( 1632 MycDeliveryTargetStatus::Unknown, 1633 attempts_remaining 1634 .then(|| next_attempt_time(job, attempt.number, retry_jitter, observed_at)) 1635 .transpose()?, 1636 ), 1637 MycDeliveryAttemptStatus::Failed => (MycDeliveryTargetStatus::Exhausted, None), 1638 MycDeliveryAttemptStatus::Leased | MycDeliveryAttemptStatus::Submitted => { 1639 return Err(DeliveryOperationError::Binding); 1640 } 1641 }; 1642 let result = sqlx::query(RESOLVE_TARGET_SQL) 1643 .bind(target_status.as_str()) 1644 .bind(next_attempt.map(MycDeliveryTimeUnixMs::sqlite_value)) 1645 .bind(observed_at.sqlite_value()) 1646 .bind(job.id.as_bytes().as_slice()) 1647 .bind(i64::from(target.index)) 1648 .bind(attempt.id.as_bytes().as_slice()) 1649 .bind(match prior { 1650 MycDeliveryAttemptStatus::Leased => "leased", 1651 MycDeliveryAttemptStatus::Submitted => "submitted", 1652 _ => return Err(DeliveryOperationError::Binding), 1653 }) 1654 .execute(&mut *transaction) 1655 .await 1656 .map_err(|_| DeliveryOperationError::Storage)?; 1657 require_one(result.rows_affected()) 1658 } 1659 1660 pub(crate) async fn finalize_job_if_terminal( 1661 transaction: &mut ServiceSqliteTransaction<'_>, 1662 job_id: MycDeliveryJobId, 1663 observed_at: MycDeliveryTimeUnixMs, 1664 ) -> Result<MycDeliveryJobRecord, DeliveryOperationError> { 1665 let job = read_job(transaction, job_id) 1666 .await? 1667 .ok_or(DeliveryOperationError::Binding)?; 1668 let delivered_required = job 1669 .targets 1670 .iter() 1671 .filter(|target| target.required && target.status == MycDeliveryTargetStatus::Delivered) 1672 .count(); 1673 let possible_required = job 1674 .targets 1675 .iter() 1676 .filter(|target| { 1677 target.required 1678 && target.status != MycDeliveryTargetStatus::Exhausted 1679 && !(target.status == MycDeliveryTargetStatus::Unknown 1680 && target.attempt_count >= job.max_attempts) 1681 }) 1682 .count(); 1683 let required = usize::try_from(job.required_acknowledgements) 1684 .map_err(|_| DeliveryOperationError::Binding)?; 1685 let delivered = match job.policy_mode { 1686 MycDeliveryPolicyMode::AtLeastOneRequired => delivered_required >= required, 1687 MycDeliveryPolicyMode::AllRequired | MycDeliveryPolicyMode::RequiredQuorum => { 1688 delivered_required >= required 1689 } 1690 }; 1691 let possible = match job.policy_mode { 1692 MycDeliveryPolicyMode::AtLeastOneRequired => possible_required >= required, 1693 MycDeliveryPolicyMode::AllRequired | MycDeliveryPolicyMode::RequiredQuorum => { 1694 possible_required >= required 1695 } 1696 }; 1697 if delivered || !possible { 1698 let terminal = if delivered { 1699 MycDeliveryJobStatus::Delivered 1700 } else if job.targets.iter().any(|target| { 1701 target.required 1702 && target.status == MycDeliveryTargetStatus::Unknown 1703 && target.attempt_count >= job.max_attempts 1704 }) { 1705 MycDeliveryJobStatus::Unknown 1706 } else { 1707 MycDeliveryJobStatus::Failed 1708 }; 1709 let result = sqlx::query(FINALIZE_JOB_SQL) 1710 .bind(terminal.as_str()) 1711 .bind(observed_at.sqlite_value()) 1712 .bind(observed_at.sqlite_value()) 1713 .bind(job_id.as_bytes().as_slice()) 1714 .execute(&mut *transaction) 1715 .await 1716 .map_err(|_| DeliveryOperationError::Storage)?; 1717 if result.rows_affected() > 1 { 1718 return Err(DeliveryOperationError::Storage); 1719 } 1720 } 1721 read_job(transaction, job_id) 1722 .await? 1723 .ok_or(DeliveryOperationError::Binding) 1724 } 1725 1726 fn next_attempt_time( 1727 job: &MycDeliveryJobRecord, 1728 attempt_number: u32, 1729 jitter: MycDeliveryRetryJitter, 1730 observed_at: MycDeliveryTimeUnixMs, 1731 ) -> Result<MycDeliveryTimeUnixMs, DeliveryOperationError> { 1732 let exponent = attempt_number.saturating_sub(1).min(31); 1733 let factor = 1_u64.checked_shl(exponent).unwrap_or(u64::MAX); 1734 let maximum_delay = job 1735 .initial_backoff_ms 1736 .saturating_mul(factor) 1737 .min(job.maximum_backoff_ms); 1738 if jitter.get() > maximum_delay { 1739 return Err(DeliveryOperationError::Binding); 1740 } 1741 let value = observed_at 1742 .get() 1743 .checked_add(jitter.get()) 1744 .ok_or(DeliveryOperationError::Binding)?; 1745 MycDeliveryTimeUnixMs::new(value).map_err(|_| DeliveryOperationError::Binding) 1746 } 1747 1748 async fn read_job_by_source( 1749 transaction: &mut ServiceSqliteTransaction<'_>, 1750 source: MycDeliverySource, 1751 ) -> Result<Option<MycDeliveryJobRecord>, DeliveryOperationError> { 1752 let rows = sqlx::query(READ_JOB_BY_SOURCE_SQL) 1753 .bind(source.kind().as_str()) 1754 .bind(source.id().as_slice()) 1755 .fetch_all(&mut *transaction) 1756 .await 1757 .map_err(|_| DeliveryOperationError::Storage)?; 1758 let job = read_job_rows(transaction, rows).await?; 1759 match job { 1760 Some(job) if job.source == source => Ok(Some(job)), 1761 Some(_) => Err(DeliveryOperationError::Binding), 1762 None => Ok(None), 1763 } 1764 } 1765 1766 pub(crate) async fn read_job( 1767 transaction: &mut ServiceSqliteTransaction<'_>, 1768 job_id: MycDeliveryJobId, 1769 ) -> Result<Option<MycDeliveryJobRecord>, DeliveryOperationError> { 1770 let rows = sqlx::query(READ_JOB_SQL) 1771 .bind(job_id.as_bytes().as_slice()) 1772 .fetch_all(&mut *transaction) 1773 .await 1774 .map_err(|_| DeliveryOperationError::Storage)?; 1775 read_job_rows(transaction, rows).await 1776 } 1777 1778 async fn read_job_rows( 1779 transaction: &mut ServiceSqliteTransaction<'_>, 1780 rows: Vec<sqlx::sqlite::SqliteRow>, 1781 ) -> Result<Option<MycDeliveryJobRecord>, DeliveryOperationError> { 1782 if rows.len() > 1 { 1783 return Err(DeliveryOperationError::Binding); 1784 } 1785 let Some(row) = rows.first() else { 1786 return Ok(None); 1787 }; 1788 let id = MycDeliveryJobId(blob32(row, "job_id")?); 1789 let source_kind = MycDeliverySourceKind::parse(text(row, "source_kind")?) 1790 .ok_or(DeliveryOperationError::Binding)?; 1791 let source = MycDeliverySource { 1792 kind: source_kind, 1793 id: blob32(row, "source_id")?, 1794 }; 1795 let artifact_digest = MycDeliveryArtifactDigest(blob32(row, "artifact_sha256")?); 1796 let policy_mode = MycDeliveryPolicyMode::parse(text(row, "policy_mode")?) 1797 .ok_or(DeliveryOperationError::Binding)?; 1798 let required_acknowledgements = positive_u32(row, "required_acknowledgements")?; 1799 let max_attempts = positive_u32(row, "max_attempts")?; 1800 if max_attempts > MYC_DELIVERY_ATTEMPT_MAX_COUNT { 1801 return Err(DeliveryOperationError::Binding); 1802 } 1803 let initial_backoff_ms = positive_u64(row, "initial_backoff_ms")?; 1804 let maximum_backoff_ms = positive_u64(row, "maximum_backoff_ms")?; 1805 let attempt_deadline_ms = positive_u64(row, "attempt_deadline_ms")?; 1806 if initial_backoff_ms > maximum_backoff_ms 1807 || maximum_backoff_ms > 300_000 1808 || attempt_deadline_ms > 30_000 1809 { 1810 return Err(DeliveryOperationError::Binding); 1811 } 1812 let status = 1813 MycDeliveryJobStatus::parse(text(row, "status")?).ok_or(DeliveryOperationError::Binding)?; 1814 let created_at = time(row, "created_at_unix_ms")?; 1815 let updated_at = time(row, "updated_at_unix_ms")?; 1816 let finalized_at = optional_time(row, "finalized_at_unix_ms", "finalized_at_type")?; 1817 let targets = read_targets(transaction, id).await?; 1818 let valid = !targets.is_empty() 1819 && targets.len() <= MYC_DELIVERY_TARGET_MAX_COUNT 1820 && targets 1821 .iter() 1822 .enumerate() 1823 .all(|(index, target)| usize::try_from(target.index) == Ok(index)) 1824 && matches!( 1825 status, 1826 MycDeliveryJobStatus::Delivered 1827 | MycDeliveryJobStatus::Failed 1828 | MycDeliveryJobStatus::Unknown 1829 ) == finalized_at.is_some() 1830 && targets.iter().all(|target| { 1831 target.status != MycDeliveryTargetStatus::Unknown 1832 || target.next_attempt_at.is_some() 1833 || target.attempt_count == max_attempts 1834 }); 1835 if !valid { 1836 return Err(DeliveryOperationError::Binding); 1837 } 1838 Ok(Some(MycDeliveryJobRecord { 1839 id, 1840 source, 1841 artifact_digest, 1842 policy_mode, 1843 required_acknowledgements, 1844 max_attempts, 1845 initial_backoff_ms, 1846 maximum_backoff_ms, 1847 attempt_deadline_ms, 1848 status, 1849 created_at, 1850 updated_at, 1851 finalized_at, 1852 targets, 1853 })) 1854 } 1855 1856 async fn read_targets( 1857 transaction: &mut ServiceSqliteTransaction<'_>, 1858 job_id: MycDeliveryJobId, 1859 ) -> Result<Box<[MycDeliveryTargetRecord]>, DeliveryOperationError> { 1860 let rows = sqlx::query(READ_TARGETS_SQL) 1861 .bind(job_id.as_bytes().as_slice()) 1862 .fetch_all(&mut *transaction) 1863 .await 1864 .map_err(|_| DeliveryOperationError::Storage)?; 1865 if rows.len() > MYC_DELIVERY_TARGET_MAX_COUNT { 1866 return Err(DeliveryOperationError::Binding); 1867 } 1868 rows.iter() 1869 .map(parse_target) 1870 .collect::<Result<Vec<_>, _>>() 1871 .map(Vec::into_boxed_slice) 1872 } 1873 1874 fn parse_target( 1875 row: &sqlx::sqlite::SqliteRow, 1876 ) -> Result<MycDeliveryTargetRecord, DeliveryOperationError> { 1877 let index = nonnegative_u32(row, "target_index")?; 1878 let relay_id = MycDeliveryRelayId::new(text(row, "relay_id")?) 1879 .map_err(|_| DeliveryOperationError::Binding)?; 1880 let required = bool_value(row, "required")?; 1881 let attempt_count = nonnegative_u32(row, "attempt_count")?; 1882 if attempt_count > MYC_DELIVERY_ATTEMPT_MAX_COUNT { 1883 return Err(DeliveryOperationError::Binding); 1884 } 1885 let status = MycDeliveryTargetStatus::parse(text(row, "status")?) 1886 .ok_or(DeliveryOperationError::Binding)?; 1887 let active_attempt_id = optional_blob32(row, "active_attempt_id", "active_attempt_id_type")? 1888 .map(MycDeliveryAttemptId); 1889 let next_attempt_at = optional_time(row, "next_attempt_at_unix_ms", "next_attempt_at_type")?; 1890 let updated_at = time(row, "updated_at_unix_ms")?; 1891 let active = matches!( 1892 status, 1893 MycDeliveryTargetStatus::Leased | MycDeliveryTargetStatus::Submitted 1894 ); 1895 let retryable = status == MycDeliveryTargetStatus::Retryable; 1896 if active != active_attempt_id.is_some() 1897 || (retryable && next_attempt_at.is_none()) 1898 || (!retryable && status != MycDeliveryTargetStatus::Unknown && next_attempt_at.is_some()) 1899 { 1900 return Err(DeliveryOperationError::Binding); 1901 } 1902 Ok(MycDeliveryTargetRecord { 1903 index, 1904 relay_id, 1905 required, 1906 attempt_count, 1907 status, 1908 active_attempt_id, 1909 next_attempt_at, 1910 updated_at, 1911 }) 1912 } 1913 1914 pub(crate) async fn read_attempts( 1915 transaction: &mut ServiceSqliteTransaction<'_>, 1916 job_id: MycDeliveryJobId, 1917 target_index: u32, 1918 ) -> Result<Box<[MycDeliveryAttemptRecord]>, DeliveryOperationError> { 1919 let rows = sqlx::query(READ_ATTEMPTS_SQL) 1920 .bind(job_id.as_bytes().as_slice()) 1921 .bind(i64::from(target_index)) 1922 .fetch_all(&mut *transaction) 1923 .await 1924 .map_err(|_| DeliveryOperationError::Storage)?; 1925 if rows.len() > usize::try_from(MYC_DELIVERY_ATTEMPT_MAX_COUNT).unwrap_or(usize::MAX) { 1926 return Err(DeliveryOperationError::Binding); 1927 } 1928 rows.iter() 1929 .map(parse_attempt) 1930 .collect::<Result<Vec<_>, _>>() 1931 .map(Vec::into_boxed_slice) 1932 } 1933 1934 async fn read_attempt( 1935 transaction: &mut ServiceSqliteTransaction<'_>, 1936 job_id: MycDeliveryJobId, 1937 target_index: u32, 1938 attempt_id: MycDeliveryAttemptId, 1939 ) -> Result<Option<MycDeliveryAttemptRecord>, DeliveryOperationError> { 1940 let rows = sqlx::query(READ_ATTEMPT_SQL) 1941 .bind(job_id.as_bytes().as_slice()) 1942 .bind(i64::from(target_index)) 1943 .bind(attempt_id.as_bytes().as_slice()) 1944 .fetch_all(&mut *transaction) 1945 .await 1946 .map_err(|_| DeliveryOperationError::Storage)?; 1947 one_attempt(rows) 1948 } 1949 1950 async fn read_attempt_by_nonce( 1951 transaction: &mut ServiceSqliteTransaction<'_>, 1952 job_id: MycDeliveryJobId, 1953 target_index: u32, 1954 nonce: &MycDeliveryAttemptNonce, 1955 ) -> Result<Option<MycDeliveryAttemptRecord>, DeliveryOperationError> { 1956 let rows = sqlx::query(READ_ATTEMPT_BY_NONCE_SQL) 1957 .bind(job_id.as_bytes().as_slice()) 1958 .bind(i64::from(target_index)) 1959 .bind(nonce.0.as_slice()) 1960 .fetch_all(&mut *transaction) 1961 .await 1962 .map_err(|_| DeliveryOperationError::Storage)?; 1963 one_attempt(rows) 1964 } 1965 1966 fn one_attempt( 1967 rows: Vec<sqlx::sqlite::SqliteRow>, 1968 ) -> Result<Option<MycDeliveryAttemptRecord>, DeliveryOperationError> { 1969 if rows.len() > 1 { 1970 return Err(DeliveryOperationError::Binding); 1971 } 1972 rows.first().map(parse_attempt).transpose() 1973 } 1974 1975 fn parse_attempt( 1976 row: &sqlx::sqlite::SqliteRow, 1977 ) -> Result<MycDeliveryAttemptRecord, DeliveryOperationError> { 1978 let id = MycDeliveryAttemptId(blob32(row, "attempt_id")?); 1979 let number = positive_u32(row, "attempt_number")?; 1980 if number > MYC_DELIVERY_ATTEMPT_MAX_COUNT { 1981 return Err(DeliveryOperationError::Binding); 1982 } 1983 let _nonce = blob32(row, "attempt_nonce")?; 1984 let status = MycDeliveryAttemptStatus::parse(text(row, "status")?) 1985 .ok_or(DeliveryOperationError::Binding)?; 1986 let leased_at = time(row, "leased_at_unix_ms")?; 1987 let lease_expires_at = time(row, "lease_expires_at_unix_ms")?; 1988 let submitted_at = optional_time(row, "submitted_at_unix_ms", "submitted_at_type")?; 1989 let resolved_at = optional_time(row, "resolved_at_unix_ms", "resolved_at_type")?; 1990 let reason = optional_reason(row)?; 1991 let valid = lease_expires_at > leased_at 1992 && match status { 1993 MycDeliveryAttemptStatus::Leased => { 1994 submitted_at.is_none() && resolved_at.is_none() && reason.is_none() 1995 } 1996 MycDeliveryAttemptStatus::Submitted => { 1997 submitted_at.is_some() && resolved_at.is_none() && reason.is_none() 1998 } 1999 MycDeliveryAttemptStatus::Delivered => { 2000 submitted_at.is_some() && resolved_at.is_some() && reason == Some("accepted") 2001 } 2002 MycDeliveryAttemptStatus::Failed => resolved_at.is_some() && reason.is_some(), 2003 MycDeliveryAttemptStatus::Unknown => { 2004 submitted_at.is_some() 2005 && resolved_at.is_some() 2006 && reason == Some("acknowledgement_lost") 2007 } 2008 }; 2009 valid 2010 .then_some(MycDeliveryAttemptRecord { 2011 id, 2012 number, 2013 status, 2014 leased_at, 2015 lease_expires_at, 2016 submitted_at, 2017 resolved_at, 2018 reason, 2019 }) 2020 .ok_or(DeliveryOperationError::Binding) 2021 } 2022 2023 fn optional_reason( 2024 row: &sqlx::sqlite::SqliteRow, 2025 ) -> Result<Option<&'static str>, DeliveryOperationError> { 2026 let kind = row 2027 .try_get::<&str, _>("reason_code_type") 2028 .map_err(|_| DeliveryOperationError::Binding)?; 2029 if kind == "null" { 2030 return Ok(None); 2031 } 2032 if kind != "text" { 2033 return Err(DeliveryOperationError::Binding); 2034 } 2035 match row 2036 .try_get::<Option<&str>, _>("reason_code") 2037 .map_err(|_| DeliveryOperationError::Binding)? 2038 .ok_or(DeliveryOperationError::Binding)? 2039 { 2040 "accepted" => Ok(Some("accepted")), 2041 "relay_rejected" => Ok(Some("relay_rejected")), 2042 "transport_failed" => Ok(Some("transport_failed")), 2043 "lease_expired_before_submit" => Ok(Some("lease_expired_before_submit")), 2044 "acknowledgement_lost" => Ok(Some("acknowledgement_lost")), 2045 _ => Err(DeliveryOperationError::Binding), 2046 } 2047 } 2048 2049 fn target_by_relay<'a>( 2050 job: &'a MycDeliveryJobRecord, 2051 relay_id: &MycDeliveryRelayId, 2052 ) -> Result<&'a MycDeliveryTargetRecord, DeliveryOperationError> { 2053 job.targets 2054 .iter() 2055 .find(|target| target.relay_id == *relay_id) 2056 .ok_or(DeliveryOperationError::Binding) 2057 } 2058 2059 fn exact_job( 2060 job: &MycDeliveryJobRecord, 2061 source: MycDeliverySource, 2062 artifact_digest: MycDeliveryArtifactDigest, 2063 created_at: MycDeliveryTimeUnixMs, 2064 policy: &MycDeliveryPolicies, 2065 ) -> bool { 2066 job.source == source 2067 && job.artifact_digest == artifact_digest 2068 && job.policy_mode == policy.mode 2069 && job.required_acknowledgements == policy.required_acknowledgements 2070 && job.max_attempts == policy.max_attempts 2071 && job.initial_backoff_ms == policy.initial_backoff_ms 2072 && job.maximum_backoff_ms == policy.maximum_backoff_ms 2073 && job.attempt_deadline_ms == policy.attempt_deadline_ms 2074 && job.created_at == created_at 2075 && job.targets.len() == policy.targets.len() 2076 && job 2077 .targets 2078 .iter() 2079 .zip(policy.targets.iter()) 2080 .all(|(actual, expected)| { 2081 actual.relay_id == expected.relay_id && actual.required == expected.required 2082 }) 2083 } 2084 2085 fn derive_job_id( 2086 source: MycDeliverySource, 2087 artifact: MycDeliveryArtifactDigest, 2088 ) -> MycDeliveryJobId { 2089 let mut hasher = Sha256::new(); 2090 hasher.update(JOB_ID_DOMAIN); 2091 if source.kind() == MycDeliverySourceKind::DiscoveryHandler { 2092 hasher.update(b"discovery_handler\0"); 2093 } 2094 hasher.update(source.id()); 2095 hasher.update(artifact.as_bytes()); 2096 MycDeliveryJobId(hasher.finalize().into()) 2097 } 2098 2099 fn derive_attempt_id( 2100 job_id: MycDeliveryJobId, 2101 target_index: u32, 2102 attempt_number: u32, 2103 nonce: &MycDeliveryAttemptNonce, 2104 ) -> MycDeliveryAttemptId { 2105 let mut hasher = Sha256::new(); 2106 hasher.update(ATTEMPT_ID_DOMAIN); 2107 hasher.update(job_id.as_bytes()); 2108 hasher.update(target_index.to_be_bytes()); 2109 hasher.update(attempt_number.to_be_bytes()); 2110 hasher.update(nonce.0); 2111 MycDeliveryAttemptId(hasher.finalize().into()) 2112 } 2113 2114 fn text<'a>( 2115 row: &'a sqlx::sqlite::SqliteRow, 2116 column: &str, 2117 ) -> Result<&'a str, DeliveryOperationError> { 2118 row.try_get::<Option<&str>, _>(column) 2119 .map_err(|_| DeliveryOperationError::Binding)? 2120 .ok_or(DeliveryOperationError::Binding) 2121 } 2122 2123 fn blob32(row: &sqlx::sqlite::SqliteRow, column: &str) -> Result<[u8; 32], DeliveryOperationError> { 2124 row.try_get::<Option<Vec<u8>>, _>(column) 2125 .map_err(|_| DeliveryOperationError::Binding)? 2126 .ok_or(DeliveryOperationError::Binding)? 2127 .try_into() 2128 .map_err(|_| DeliveryOperationError::Binding) 2129 } 2130 2131 fn optional_blob32( 2132 row: &sqlx::sqlite::SqliteRow, 2133 column: &str, 2134 type_column: &str, 2135 ) -> Result<Option<[u8; 32]>, DeliveryOperationError> { 2136 match row 2137 .try_get::<&str, _>(type_column) 2138 .map_err(|_| DeliveryOperationError::Binding)? 2139 { 2140 "null" => Ok(None), 2141 "blob" => blob32(row, column).map(Some), 2142 _ => Err(DeliveryOperationError::Binding), 2143 } 2144 } 2145 2146 fn positive_u32( 2147 row: &sqlx::sqlite::SqliteRow, 2148 column: &str, 2149 ) -> Result<u32, DeliveryOperationError> { 2150 nonnegative_u32(row, column).and_then(|value| { 2151 (value != 0) 2152 .then_some(value) 2153 .ok_or(DeliveryOperationError::Binding) 2154 }) 2155 } 2156 2157 fn nonnegative_u32( 2158 row: &sqlx::sqlite::SqliteRow, 2159 column: &str, 2160 ) -> Result<u32, DeliveryOperationError> { 2161 let value = row 2162 .try_get::<i64, _>(column) 2163 .map_err(|_| DeliveryOperationError::Binding)?; 2164 u32::try_from(value).map_err(|_| DeliveryOperationError::Binding) 2165 } 2166 2167 fn positive_u64( 2168 row: &sqlx::sqlite::SqliteRow, 2169 column: &str, 2170 ) -> Result<u64, DeliveryOperationError> { 2171 let value = row 2172 .try_get::<i64, _>(column) 2173 .map_err(|_| DeliveryOperationError::Binding)?; 2174 u64::try_from(value) 2175 .ok() 2176 .filter(|value| *value != 0) 2177 .ok_or(DeliveryOperationError::Binding) 2178 } 2179 2180 fn bool_value(row: &sqlx::sqlite::SqliteRow, column: &str) -> Result<bool, DeliveryOperationError> { 2181 match row 2182 .try_get::<i64, _>(column) 2183 .map_err(|_| DeliveryOperationError::Binding)? 2184 { 2185 0 => Ok(false), 2186 1 => Ok(true), 2187 _ => Err(DeliveryOperationError::Binding), 2188 } 2189 } 2190 2191 fn time( 2192 row: &sqlx::sqlite::SqliteRow, 2193 column: &str, 2194 ) -> Result<MycDeliveryTimeUnixMs, DeliveryOperationError> { 2195 let value = positive_u64(row, column)?; 2196 MycDeliveryTimeUnixMs::new(value).map_err(|_| DeliveryOperationError::Binding) 2197 } 2198 2199 fn optional_time( 2200 row: &sqlx::sqlite::SqliteRow, 2201 column: &str, 2202 type_column: &str, 2203 ) -> Result<Option<MycDeliveryTimeUnixMs>, DeliveryOperationError> { 2204 match row 2205 .try_get::<&str, _>(type_column) 2206 .map_err(|_| DeliveryOperationError::Binding)? 2207 { 2208 "null" => Ok(None), 2209 "integer" => time(row, column).map(Some), 2210 _ => Err(DeliveryOperationError::Binding), 2211 } 2212 } 2213 2214 fn require_one(rows: u64) -> Result<(), DeliveryOperationError> { 2215 (rows == 1) 2216 .then_some(()) 2217 .ok_or(DeliveryOperationError::Storage) 2218 } 2219 2220 fn map_transaction_error( 2221 error: ServiceSqliteTransactionError<DeliveryOperationError>, 2222 ) -> MycStateRepositoryError { 2223 if error.kind() == ServiceSqliteTransactionErrorKind::CommitOutcomeUnknown { 2224 return MycStateRepositoryError::new(MycStateRepositoryErrorKind::CommitOutcomeUnknown); 2225 } 2226 let kind = match error.operation_error() { 2227 Some(DeliveryOperationError::Binding) => MycStateRepositoryErrorKind::Binding, 2228 Some(DeliveryOperationError::Storage) | None => MycStateRepositoryErrorKind::Transaction, 2229 }; 2230 MycStateRepositoryError::new(kind) 2231 }