mod.rs (47395B)
1 use crate::SqliteStorage; 2 use crate::backend::map_backend; 3 use radroots_event_codec::Codec; 4 use radroots_storage::{ 5 Error, Outbox, 6 outbox::{ 7 BoxFuture, ClaimOutboxItems, ClaimedOutboxItem, DeliveryAttempt, DeliveryAttemptEvidence, 8 DeliveryOutcome, DeliveryOutcomeKind, DeliveryPayload, DeliveryPlanDigest, DeliveryRequest, 9 EnqueueDisposition, EnqueueOutboxItem, EnqueueReceipt, LeaseId, LeaseOwner, OutboxItemId, 10 OutboxLease, OutboxRecord, OutboxRevision, OutboxStage, OutboxStatus, Retryability, 11 SatisfactionClass, SatisfactionPolicy, SatisfactionResult, TARGET_SET_MAX_ITEMS, Target, 12 TargetDeliveryEvidence, TargetFingerprint, TargetLabel, TargetPolicy, TargetScope, 13 TargetSet, TransportId, 14 }, 15 }; 16 use sqlx::{Row, Sqlite, SqliteConnection}; 17 18 #[cfg_attr(coverage_nightly, coverage(off))] 19 impl Outbox for SqliteStorage { 20 fn enqueue(&self, item: EnqueueOutboxItem) -> BoxFuture<'_, Result<EnqueueReceipt, Error>> { 21 Box::pin(async move { 22 self.require_outbox_writer()?; 23 let mut transaction = self 24 .pool() 25 .begin_with("BEGIN IMMEDIATE") 26 .await 27 .map_err(map_backend)?; 28 let receipt = enqueue_transaction(&mut transaction, item).await?; 29 transaction.commit().await.map_err(map_backend)?; 30 Ok(receipt) 31 }) 32 } 33 34 fn item(&self, item_id: OutboxItemId) -> BoxFuture<'_, Result<Option<OutboxRecord>, Error>> { 35 Box::pin(async move { 36 let mut connection = self.pool().acquire().await.map_err(map_backend)?; 37 load_record(&mut connection, item_id).await 38 }) 39 } 40 41 fn claim( 42 &self, 43 request: ClaimOutboxItems, 44 ) -> BoxFuture<'_, Result<Vec<ClaimedOutboxItem>, Error>> { 45 Box::pin(async move { 46 self.require_outbox_writer()?; 47 let mut transaction = self 48 .pool() 49 .begin_with("BEGIN IMMEDIATE") 50 .await 51 .map_err(map_backend)?; 52 let rows = sqlx::query( 53 "SELECT item_id FROM radroots_runtime_outbox_items 54 WHERE stage IN ('pending', 'leased', 'retryable') 55 AND (retry_not_before_unix_ms IS NULL OR retry_not_before_unix_ms <= ?) 56 AND (stage <> 'leased' OR lease_expires_at_unix_ms <= ?) 57 ORDER BY created_at_unix_ms, item_id LIMIT ?", 58 ) 59 .bind(i64_from_u64(request.now_unix_ms())?) 60 .bind(i64_from_u64(request.now_unix_ms())?) 61 .bind(i64::from(request.limit())) 62 .fetch_all(&mut *transaction) 63 .await 64 .map_err(map_backend)?; 65 let item_ids = rows 66 .iter() 67 .map(|row| { 68 OutboxItemId::new(array( 69 row.try_get::<Vec<u8>, _>("item_id").map_err(map_corrupt)?, 70 )?) 71 .map_err(|_| Error::CorruptOutboxRecord) 72 }) 73 .collect::<Result<Vec<_>, _>>()?; 74 75 let mut claimed = Vec::with_capacity(item_ids.len()); 76 for item_id in item_ids { 77 let mut record = load_record(&mut transaction, item_id) 78 .await? 79 .ok_or(Error::CorruptOutboxRecord)?; 80 let prior_revision = record.revision(); 81 let lease = OutboxLease::new( 82 request.lease_id_for(item_id), 83 request.owner().clone(), 84 request.now_unix_ms(), 85 request.lease_expires_at_unix_ms(), 86 )?; 87 record.claim(lease.clone())?; 88 update_record(&mut transaction, &record, prior_revision).await?; 89 claimed.push(ClaimedOutboxItem::new(record, lease)); 90 } 91 transaction.commit().await.map_err(map_backend)?; 92 Ok(claimed) 93 }) 94 } 95 96 fn record_attempt( 97 &self, 98 evidence: DeliveryAttemptEvidence, 99 ) -> BoxFuture<'_, Result<OutboxRecord, Error>> { 100 Box::pin(async move { 101 self.require_outbox_writer()?; 102 let mut transaction = self 103 .pool() 104 .begin_with("BEGIN IMMEDIATE") 105 .await 106 .map_err(map_backend)?; 107 let record = record_attempt_transaction(&mut transaction, evidence).await?; 108 transaction.commit().await.map_err(map_backend)?; 109 Ok(record) 110 }) 111 } 112 113 fn release( 114 &self, 115 item_id: OutboxItemId, 116 lease_id: LeaseId, 117 expected_revision: OutboxRevision, 118 released_at_unix_ms: u64, 119 retry_not_before_unix_ms: Option<u64>, 120 ) -> BoxFuture<'_, Result<OutboxRecord, Error>> { 121 Box::pin(async move { 122 self.require_outbox_writer()?; 123 let mut transaction = self 124 .pool() 125 .begin_with("BEGIN IMMEDIATE") 126 .await 127 .map_err(map_backend)?; 128 let mut record = load_record(&mut transaction, item_id) 129 .await? 130 .ok_or(Error::OutboxItemNotFound)?; 131 let prior_revision = record.revision(); 132 record.release( 133 lease_id, 134 expected_revision, 135 released_at_unix_ms, 136 retry_not_before_unix_ms, 137 )?; 138 update_record(&mut transaction, &record, prior_revision).await?; 139 transaction.commit().await.map_err(map_backend)?; 140 Ok(record) 141 }) 142 } 143 144 fn status(&self) -> BoxFuture<'_, Result<OutboxStatus, Error>> { 145 Box::pin(async move { 146 let row = sqlx::query( 147 "SELECT 148 COALESCE(SUM(CASE WHEN stage = 'pending' THEN 1 ELSE 0 END), 0) AS pending, 149 COALESCE(SUM(CASE WHEN stage = 'leased' THEN 1 ELSE 0 END), 0) AS leased, 150 COALESCE(SUM(CASE WHEN stage = 'retryable' THEN 1 ELSE 0 END), 0) AS retryable, 151 COALESCE(SUM(CASE WHEN stage = 'satisfied' THEN 1 ELSE 0 END), 0) AS satisfied, 152 COALESCE(SUM(CASE WHEN stage = 'exhausted' THEN 1 ELSE 0 END), 0) AS exhausted 153 FROM radroots_runtime_outbox_items", 154 ) 155 .fetch_one(self.pool()) 156 .await 157 .map_err(map_backend)?; 158 Ok(OutboxStatus { 159 pending: count(&row, "pending")?, 160 leased: count(&row, "leased")?, 161 retryable: count(&row, "retryable")?, 162 satisfied: count(&row, "satisfied")?, 163 exhausted: count(&row, "exhausted")?, 164 }) 165 }) 166 } 167 } 168 169 #[cfg_attr(coverage_nightly, coverage(off))] 170 pub(crate) async fn enqueue_transaction( 171 transaction: &mut sqlx::Transaction<'_, Sqlite>, 172 item: EnqueueOutboxItem, 173 ) -> Result<EnqueueReceipt, Error> { 174 if let Some(record) = load_record(transaction, item.item_id()).await? { 175 if record.operation_instance_id() != item.operation_instance_id() 176 || record.plan_digest() != item.plan_digest() 177 || record.request() != item.request() 178 || record.created_at_unix_ms() != item.created_at_unix_ms() 179 { 180 return Err(Error::OutboxPlanConflict); 181 } 182 return Ok(EnqueueReceipt::new(EnqueueDisposition::Replay, record)); 183 } 184 if sqlx::query_scalar::<_, i64>( 185 "SELECT 1 FROM radroots_runtime_outbox_items WHERE operation_instance_id = ?", 186 ) 187 .bind(item.operation_instance_id().as_bytes().as_slice()) 188 .fetch_optional(&mut **transaction) 189 .await 190 .map_err(map_backend)? 191 .is_some() 192 { 193 return Err(Error::OutboxPlanConflict); 194 } 195 if sqlx::query_scalar::<_, i64>( 196 "SELECT 1 FROM radroots_runtime_journal_operations WHERE instance_id = ?", 197 ) 198 .bind(item.operation_instance_id().as_bytes().as_slice()) 199 .fetch_optional(&mut **transaction) 200 .await 201 .map_err(map_backend)? 202 .is_none() 203 { 204 return Err(Error::OperationNotFound); 205 } 206 207 let record = item.into_record(); 208 insert_record(transaction, &record).await?; 209 Ok(EnqueueReceipt::new(EnqueueDisposition::Created, record)) 210 } 211 212 #[cfg_attr(coverage_nightly, coverage(off))] 213 pub(crate) async fn record_attempt_transaction( 214 transaction: &mut sqlx::Transaction<'_, Sqlite>, 215 evidence: DeliveryAttemptEvidence, 216 ) -> Result<OutboxRecord, Error> { 217 let mut record = load_record(transaction, evidence.item_id()) 218 .await? 219 .ok_or(Error::OutboxItemNotFound)?; 220 let prior_revision = record.revision(); 221 let receipt = evidence.receipt().clone(); 222 let attempt = evidence.attempt(); 223 let recorded_at = evidence.recorded_at_unix_ms(); 224 record.record_attempt(evidence)?; 225 update_record(transaction, &record, prior_revision).await?; 226 for target_receipt in receipt.target_receipts() { 227 let outcome = encode_outcome(target_receipt.outcome())?; 228 sqlx::query( 229 "INSERT INTO radroots_runtime_delivery_evidence ( 230 item_id, target_fingerprint, attempt, attempted, outcome, 231 retryability, recorded_at_unix_ms 232 ) VALUES (?, ?, ?, ?, ?, ?, ?)", 233 ) 234 .bind(record.item_id().as_bytes().as_slice()) 235 .bind(target_receipt.target().fingerprint().as_str().as_bytes()) 236 .bind(i64::from(attempt.get())) 237 .bind(i64::from(target_receipt.was_attempted())) 238 .bind(outcome) 239 .bind(retryability_name(target_receipt.outcome().retryability())) 240 .bind(i64_from_u64(recorded_at)?) 241 .execute(&mut **transaction) 242 .await 243 .map_err(map_backend)?; 244 } 245 Ok(record) 246 } 247 248 impl SqliteStorage { 249 fn require_outbox_writer(&self) -> Result<(), Error> { 250 if self.event_mode() == radroots_storage::status::EventStoreMode::ReadOnly { 251 return Err(Error::BackendUnavailable); 252 } 253 Ok(()) 254 } 255 } 256 257 #[cfg_attr(coverage_nightly, coverage(off))] 258 async fn insert_record( 259 transaction: &mut sqlx::Transaction<'_, Sqlite>, 260 record: &OutboxRecord, 261 ) -> Result<(), Error> { 262 let request = encode_request(record.request())?; 263 sqlx::query( 264 "INSERT INTO radroots_runtime_outbox_items ( 265 item_id, operation_instance_id, plan_digest, delivery_request, revision, stage, 266 satisfaction, created_at_unix_ms, updated_at_unix_ms 267 ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)", 268 ) 269 .bind(record.item_id().as_bytes().as_slice()) 270 .bind(record.operation_instance_id().as_bytes().as_slice()) 271 .bind(record.plan_digest().as_bytes().as_slice()) 272 .bind(request) 273 .bind(i64_from_u64(record.revision().get())?) 274 .bind(stage_name(record.stage())) 275 .bind(satisfaction_name(record.satisfaction())) 276 .bind(i64_from_u64(record.created_at_unix_ms())?) 277 .bind(i64_from_u64(record.updated_at_unix_ms())?) 278 .execute(&mut **transaction) 279 .await 280 .map_err(map_backend)?; 281 282 for (ordinal, target) in record.request().target_set().targets().iter().enumerate() { 283 sqlx::query( 284 "INSERT INTO radroots_runtime_outbox_targets ( 285 item_id, target_fingerprint, target_request, ordinal 286 ) VALUES (?, ?, ?, ?)", 287 ) 288 .bind(record.item_id().as_bytes().as_slice()) 289 .bind(target.fingerprint().as_str().as_bytes()) 290 .bind(encode_target(target)?) 291 .bind(i64::try_from(ordinal).map_err(|_| Error::CorruptOutboxRecord)?) 292 .execute(&mut **transaction) 293 .await 294 .map_err(map_backend)?; 295 } 296 Ok(()) 297 } 298 299 #[cfg_attr(coverage_nightly, coverage(off))] 300 async fn update_record( 301 transaction: &mut sqlx::Transaction<'_, Sqlite>, 302 record: &OutboxRecord, 303 prior_revision: OutboxRevision, 304 ) -> Result<(), Error> { 305 let lease = record.lease(); 306 let result = sqlx::query( 307 "UPDATE radroots_runtime_outbox_items SET 308 revision = ?, stage = ?, lease_id = ?, lease_owner = ?, 309 lease_acquired_at_unix_ms = ?, lease_expires_at_unix_ms = ?, last_attempt = ?, 310 satisfaction = ?, retry_not_before_unix_ms = ?, updated_at_unix_ms = ? 311 WHERE item_id = ? AND revision = ?", 312 ) 313 .bind(i64_from_u64(record.revision().get())?) 314 .bind(stage_name(record.stage())) 315 .bind(lease.map(|lease| lease.id().as_bytes().to_vec())) 316 .bind(lease.map(|lease| lease.owner().as_str())) 317 .bind( 318 lease 319 .map(|lease| i64_from_u64(lease.acquired_at_unix_ms())) 320 .transpose()?, 321 ) 322 .bind( 323 lease 324 .map(|lease| i64_from_u64(lease.expires_at_unix_ms())) 325 .transpose()?, 326 ) 327 .bind( 328 record 329 .last_attempt() 330 .map(|attempt| i64::from(attempt.get())), 331 ) 332 .bind(satisfaction_name(record.satisfaction())) 333 .bind( 334 record 335 .retry_not_before_unix_ms() 336 .map(i64_from_u64) 337 .transpose()?, 338 ) 339 .bind(i64_from_u64(record.updated_at_unix_ms())?) 340 .bind(record.item_id().as_bytes().as_slice()) 341 .bind(i64_from_u64(prior_revision.get())?) 342 .execute(&mut **transaction) 343 .await 344 .map_err(map_backend)?; 345 if result.rows_affected() != 1 { 346 return Err(Error::OutboxRevisionConflict); 347 } 348 Ok(()) 349 } 350 351 #[cfg_attr(coverage_nightly, coverage(off))] 352 pub(crate) async fn load_record( 353 connection: &mut SqliteConnection, 354 item_id: OutboxItemId, 355 ) -> Result<Option<OutboxRecord>, Error> { 356 let Some(row) = sqlx::query("SELECT * FROM radroots_runtime_outbox_items WHERE item_id = ?") 357 .bind(item_id.as_bytes().as_slice()) 358 .fetch_optional(&mut *connection) 359 .await 360 .map_err(map_backend)? 361 else { 362 return Ok(None); 363 }; 364 let request = decode_request( 365 row.try_get::<Vec<u8>, _>("delivery_request") 366 .map_err(map_corrupt)? 367 .as_slice(), 368 )?; 369 validate_targets(connection, item_id, &request).await?; 370 let evidence = load_evidence(connection, item_id).await?; 371 let operation_instance_id = radroots_storage::journal::OperationInstanceId::new(array( 372 row.try_get::<Vec<u8>, _>("operation_instance_id") 373 .map_err(map_corrupt)?, 374 )?) 375 .map_err(|_| Error::CorruptOutboxRecord)?; 376 let enqueue = EnqueueOutboxItem::new( 377 item_id, 378 operation_instance_id, 379 DeliveryPlanDigest::new(array( 380 row.try_get::<Vec<u8>, _>("plan_digest") 381 .map_err(map_corrupt)?, 382 )?), 383 request, 384 u64_from_i64(row.try_get("created_at_unix_ms").map_err(map_corrupt)?)?, 385 ) 386 .map_err(|_| Error::CorruptOutboxRecord)?; 387 let lease = decode_lease(&row)?; 388 let last_attempt = row 389 .try_get::<Option<i64>, _>("last_attempt") 390 .map_err(map_corrupt)? 391 .map(|value| { 392 DeliveryAttempt::new(u32::try_from(value).map_err(|_| Error::CorruptOutboxRecord)?) 393 .map_err(|_| Error::CorruptOutboxRecord) 394 }) 395 .transpose()?; 396 OutboxRecord::from_durable_parts( 397 enqueue, 398 OutboxRevision::new(u64_from_i64(row.try_get("revision").map_err(map_corrupt)?)?) 399 .map_err(|_| Error::CorruptOutboxRecord)?, 400 stage( 401 row.try_get::<String, _>("stage") 402 .map_err(map_corrupt)? 403 .as_str(), 404 )?, 405 lease, 406 last_attempt, 407 evidence, 408 satisfaction( 409 row.try_get::<String, _>("satisfaction") 410 .map_err(map_corrupt)? 411 .as_str(), 412 )?, 413 row.try_get::<Option<i64>, _>("retry_not_before_unix_ms") 414 .map_err(map_corrupt)? 415 .map(u64_from_i64) 416 .transpose()?, 417 u64_from_i64(row.try_get("updated_at_unix_ms").map_err(map_corrupt)?)?, 418 ) 419 .map(Some) 420 } 421 422 #[cfg_attr(coverage_nightly, coverage(off))] 423 async fn validate_targets( 424 connection: &mut SqliteConnection, 425 item_id: OutboxItemId, 426 request: &DeliveryRequest, 427 ) -> Result<(), Error> { 428 let rows = sqlx::query( 429 "SELECT target_fingerprint, target_request, ordinal 430 FROM radroots_runtime_outbox_targets WHERE item_id = ? ORDER BY ordinal", 431 ) 432 .bind(item_id.as_bytes().as_slice()) 433 .fetch_all(&mut *connection) 434 .await 435 .map_err(map_backend)?; 436 if rows.len() != request.target_set().len() { 437 return Err(Error::CorruptOutboxRecord); 438 } 439 for (ordinal, (row, expected)) in rows.iter().zip(request.target_set().targets()).enumerate() { 440 let stored_ordinal = row.try_get::<i64, _>("ordinal").map_err(map_corrupt)?; 441 let fingerprint = String::from_utf8( 442 row.try_get::<Vec<u8>, _>("target_fingerprint") 443 .map_err(map_corrupt)?, 444 ) 445 .map_err(|_| Error::CorruptOutboxRecord)?; 446 let target = decode_target( 447 row.try_get::<Vec<u8>, _>("target_request") 448 .map_err(map_corrupt)? 449 .as_slice(), 450 )?; 451 if stored_ordinal != i64::try_from(ordinal).map_err(|_| Error::CorruptOutboxRecord)? 452 || fingerprint != expected.fingerprint().as_str() 453 || target != *expected 454 { 455 return Err(Error::CorruptOutboxRecord); 456 } 457 } 458 Ok(()) 459 } 460 461 #[cfg_attr(coverage_nightly, coverage(off))] 462 async fn load_evidence( 463 connection: &mut SqliteConnection, 464 item_id: OutboxItemId, 465 ) -> Result<Vec<TargetDeliveryEvidence>, Error> { 466 sqlx::query( 467 "SELECT evidence.target_fingerprint, evidence.attempt, evidence.attempted, 468 evidence.outcome, evidence.retryability, evidence.recorded_at_unix_ms 469 FROM radroots_runtime_delivery_evidence AS evidence 470 JOIN radroots_runtime_outbox_targets AS target 471 ON target.item_id = evidence.item_id 472 AND target.target_fingerprint = evidence.target_fingerprint 473 WHERE evidence.item_id = ? 474 ORDER BY evidence.attempt, target.ordinal", 475 ) 476 .bind(item_id.as_bytes().as_slice()) 477 .fetch_all(&mut *connection) 478 .await 479 .map_err(map_backend)? 480 .iter() 481 .map(|row| { 482 let target = TargetFingerprint::parse( 483 String::from_utf8( 484 row.try_get::<Vec<u8>, _>("target_fingerprint") 485 .map_err(map_corrupt)?, 486 ) 487 .map_err(|_| Error::CorruptOutboxRecord)?, 488 ) 489 .map_err(|_| Error::CorruptOutboxRecord)?; 490 let outcome = decode_outcome( 491 row.try_get::<Vec<u8>, _>("outcome") 492 .map_err(map_corrupt)? 493 .as_slice(), 494 )?; 495 let retryability = row 496 .try_get::<String, _>("retryability") 497 .map_err(map_corrupt)?; 498 if retryability != retryability_name(outcome.retryability()) { 499 return Err(Error::CorruptOutboxRecord); 500 } 501 TargetDeliveryEvidence::new( 502 target, 503 DeliveryAttempt::new( 504 u32::try_from(row.try_get::<i64, _>("attempt").map_err(map_corrupt)?) 505 .map_err(|_| Error::CorruptOutboxRecord)?, 506 ) 507 .map_err(|_| Error::CorruptOutboxRecord)?, 508 match row.try_get::<i64, _>("attempted").map_err(map_corrupt)? { 509 0 => false, 510 1 => true, 511 _ => return Err(Error::CorruptOutboxRecord), 512 }, 513 outcome, 514 u64_from_i64(row.try_get("recorded_at_unix_ms").map_err(map_corrupt)?)?, 515 ) 516 .map_err(|_| Error::CorruptOutboxRecord) 517 }) 518 .collect() 519 } 520 521 fn decode_lease(row: &sqlx::sqlite::SqliteRow) -> Result<Option<OutboxLease>, Error> { 522 let id = row 523 .try_get::<Option<Vec<u8>>, _>("lease_id") 524 .map_err(map_corrupt)?; 525 let owner = row 526 .try_get::<Option<String>, _>("lease_owner") 527 .map_err(map_corrupt)?; 528 let acquired = row 529 .try_get::<Option<i64>, _>("lease_acquired_at_unix_ms") 530 .map_err(map_corrupt)?; 531 let expires = row 532 .try_get::<Option<i64>, _>("lease_expires_at_unix_ms") 533 .map_err(map_corrupt)?; 534 match (id, owner, acquired, expires) { 535 (None, None, None, None) => Ok(None), 536 (Some(id), Some(owner), Some(acquired), Some(expires)) => OutboxLease::new( 537 LeaseId::new(array(id)?).map_err(|_| Error::CorruptOutboxRecord)?, 538 LeaseOwner::parse(owner).map_err(|_| Error::CorruptOutboxRecord)?, 539 u64_from_i64(acquired)?, 540 u64_from_i64(expires)?, 541 ) 542 .map(Some) 543 .map_err(|_| Error::CorruptOutboxRecord), 544 _ => Err(Error::CorruptOutboxRecord), 545 } 546 } 547 548 fn encode_request(value: &DeliveryRequest) -> Result<Vec<u8>, Error> { 549 let mut bytes = vec![1]; 550 put_str(&mut bytes, value.request_id().as_str())?; 551 put_blob(&mut bytes, value.payload().event().raw_json().as_bytes())?; 552 put_u16( 553 &mut bytes, 554 u16::try_from(value.target_set().len()).map_err(|_| Error::CorruptOutboxRecord)?, 555 ); 556 for target in value.target_set().targets() { 557 put_blob(&mut bytes, encode_target(target)?.as_slice())?; 558 } 559 bytes.push(match value.satisfaction().class() { 560 SatisfactionClass::Accepted => 0, 561 SatisfactionClass::Delivered => 1, 562 }); 563 let policy = value.satisfaction().targets(); 564 if policy.is_any() { 565 bytes.push(0); 566 } else if policy.is_all() { 567 bytes.push(1); 568 } else if let Some(threshold) = policy.quorum_threshold() { 569 bytes.push(2); 570 put_u16(&mut bytes, threshold); 571 } else if let Some(required) = policy.required_targets() { 572 bytes.push(3); 573 put_u16( 574 &mut bytes, 575 u16::try_from(required.len()).map_err(|_| Error::CorruptOutboxRecord)?, 576 ); 577 for target in required { 578 put_str(&mut bytes, target.as_str())?; 579 } 580 } else { 581 return Err(Error::CorruptOutboxRecord); 582 } 583 bytes.extend_from_slice(&value.deadline_unix_ms().to_be_bytes()); 584 Ok(bytes) 585 } 586 587 fn decode_request(bytes: &[u8]) -> Result<DeliveryRequest, Error> { 588 let mut cursor = Cursor::new(bytes); 589 if cursor.byte()? != 1 { 590 return Err(Error::CorruptOutboxRecord); 591 } 592 let request_id = cursor.string()?.to_owned(); 593 let raw_event = cursor.string_blob()?; 594 let event = Codec::decode_signed_event(raw_event).map_err(|_| Error::CorruptOutboxRecord)?; 595 let target_count = usize::from(cursor.u16()?); 596 if target_count == 0 || target_count > TARGET_SET_MAX_ITEMS { 597 return Err(Error::CorruptOutboxRecord); 598 } 599 let mut targets = Vec::with_capacity(target_count); 600 for _ in 0..target_count { 601 targets.push(decode_target(cursor.blob()?)?); 602 } 603 let class = match cursor.byte()? { 604 0 => SatisfactionClass::Accepted, 605 1 => SatisfactionClass::Delivered, 606 _ => return Err(Error::CorruptOutboxRecord), 607 }; 608 let target_policy = match cursor.byte()? { 609 0 => TargetPolicy::any(), 610 1 => TargetPolicy::all(), 611 2 => TargetPolicy::quorum(cursor.u16()?).map_err(|_| Error::CorruptOutboxRecord)?, 612 3 => { 613 let count = usize::from(cursor.u16()?); 614 let mut required = Vec::with_capacity(count); 615 for _ in 0..count { 616 required.push( 617 TargetFingerprint::parse(cursor.string()?) 618 .map_err(|_| Error::CorruptOutboxRecord)?, 619 ); 620 } 621 TargetPolicy::required(required).map_err(|_| Error::CorruptOutboxRecord)? 622 } 623 _ => return Err(Error::CorruptOutboxRecord), 624 }; 625 let deadline = cursor.u64()?; 626 cursor.finish()?; 627 DeliveryRequest::new( 628 request_id, 629 DeliveryPayload::new(event), 630 TargetSet::new(targets).map_err(|_| Error::CorruptOutboxRecord)?, 631 SatisfactionPolicy::new(class, target_policy), 632 deadline, 633 ) 634 .map_err(|_| Error::CorruptOutboxRecord) 635 } 636 637 fn encode_target(value: &Target) -> Result<Vec<u8>, Error> { 638 let mut bytes = vec![1]; 639 put_str(&mut bytes, value.kind().as_str())?; 640 put_str(&mut bytes, value.uri().as_str())?; 641 put_optional_str(&mut bytes, value.scope().map(TargetScope::as_str))?; 642 put_optional_str(&mut bytes, value.label().map(TargetLabel::as_str))?; 643 Ok(bytes) 644 } 645 646 fn decode_target(bytes: &[u8]) -> Result<Target, Error> { 647 let mut cursor = Cursor::new(bytes); 648 if cursor.byte()? != 1 { 649 return Err(Error::CorruptOutboxRecord); 650 } 651 let kind = TransportId::parse(cursor.string()?).map_err(|_| Error::CorruptOutboxRecord)?; 652 let uri = cursor.string()?.to_owned(); 653 let scope = cursor 654 .optional_string()? 655 .map(TargetScope::parse) 656 .transpose() 657 .map_err(|_| Error::CorruptOutboxRecord)?; 658 let label = cursor 659 .optional_string()? 660 .map(TargetLabel::parse) 661 .transpose() 662 .map_err(|_| Error::CorruptOutboxRecord)?; 663 cursor.finish()?; 664 Target::new_with_metadata(kind, uri, scope, label).map_err(|_| Error::CorruptOutboxRecord) 665 } 666 667 fn encode_outcome(value: &DeliveryOutcome) -> Result<Vec<u8>, Error> { 668 let mut bytes = vec![1]; 669 bytes.push(match value.kind() { 670 DeliveryOutcomeKind::Accepted => 0, 671 DeliveryOutcomeKind::Delivered => 1, 672 DeliveryOutcomeKind::Rejected => 2, 673 DeliveryOutcomeKind::Unavailable => 3, 674 DeliveryOutcomeKind::Failed => 4, 675 }); 676 bytes.push(retryability_byte(value.retryability())); 677 match (value.code(), value.message()) { 678 (None, None) => bytes.push(0), 679 (Some(code), Some(message)) => { 680 bytes.push(1); 681 put_str(&mut bytes, code)?; 682 put_str(&mut bytes, message)?; 683 } 684 _ => return Err(Error::CorruptOutboxRecord), 685 } 686 Ok(bytes) 687 } 688 689 fn decode_outcome(bytes: &[u8]) -> Result<DeliveryOutcome, Error> { 690 let mut cursor = Cursor::new(bytes); 691 if cursor.byte()? != 1 { 692 return Err(Error::CorruptOutboxRecord); 693 } 694 let kind = cursor.byte()?; 695 let retryability = retryability(cursor.byte()?)?; 696 let outcome = match kind { 697 0 if retryability == Retryability::NotApplicable => DeliveryOutcome::accepted(), 698 1 if retryability == Retryability::NotApplicable => DeliveryOutcome::delivered(), 699 2 if retryability == Retryability::Terminal => DeliveryOutcome::rejected(), 700 3 if retryability == Retryability::Retryable => DeliveryOutcome::unavailable(), 701 4 => DeliveryOutcome::failed(retryability).map_err(|_| Error::CorruptOutboxRecord)?, 702 _ => return Err(Error::CorruptOutboxRecord), 703 }; 704 let outcome = match cursor.byte()? { 705 0 => outcome, 706 1 => outcome 707 .with_detail(cursor.string()?, cursor.string()?) 708 .map_err(|_| Error::CorruptOutboxRecord)?, 709 _ => return Err(Error::CorruptOutboxRecord), 710 }; 711 cursor.finish()?; 712 Ok(outcome) 713 } 714 715 fn put_str(bytes: &mut Vec<u8>, value: &str) -> Result<(), Error> { 716 let length = u16::try_from(value.len()).map_err(|_| Error::CorruptOutboxRecord)?; 717 put_u16(bytes, length); 718 bytes.extend_from_slice(value.as_bytes()); 719 Ok(()) 720 } 721 722 fn put_optional_str(bytes: &mut Vec<u8>, value: Option<&str>) -> Result<(), Error> { 723 match value { 724 Some(value) => { 725 bytes.push(1); 726 put_str(bytes, value) 727 } 728 None => { 729 bytes.push(0); 730 Ok(()) 731 } 732 } 733 } 734 735 fn put_blob(bytes: &mut Vec<u8>, value: &[u8]) -> Result<(), Error> { 736 let length = u32::try_from(value.len()).map_err(|_| Error::CorruptOutboxRecord)?; 737 bytes.extend_from_slice(&length.to_be_bytes()); 738 bytes.extend_from_slice(value); 739 Ok(()) 740 } 741 742 fn put_u16(bytes: &mut Vec<u8>, value: u16) { 743 bytes.extend_from_slice(&value.to_be_bytes()); 744 } 745 746 struct Cursor<'a> { 747 bytes: &'a [u8], 748 offset: usize, 749 } 750 751 impl<'a> Cursor<'a> { 752 const fn new(bytes: &'a [u8]) -> Self { 753 Self { bytes, offset: 0 } 754 } 755 756 fn byte(&mut self) -> Result<u8, Error> { 757 let value = self 758 .bytes 759 .get(self.offset) 760 .copied() 761 .ok_or(Error::CorruptOutboxRecord)?; 762 self.offset += 1; 763 Ok(value) 764 } 765 766 fn u16(&mut self) -> Result<u16, Error> { 767 Ok(u16::from_be_bytes(self.array()?)) 768 } 769 770 fn u32(&mut self) -> Result<u32, Error> { 771 Ok(u32::from_be_bytes(self.array()?)) 772 } 773 774 fn u64(&mut self) -> Result<u64, Error> { 775 Ok(u64::from_be_bytes(self.array()?)) 776 } 777 778 fn string(&mut self) -> Result<&'a str, Error> { 779 let length = usize::from(self.u16()?); 780 core::str::from_utf8(self.take(length)?).map_err(|_| Error::CorruptOutboxRecord) 781 } 782 783 fn optional_string(&mut self) -> Result<Option<&'a str>, Error> { 784 match self.byte()? { 785 0 => Ok(None), 786 1 => self.string().map(Some), 787 _ => Err(Error::CorruptOutboxRecord), 788 } 789 } 790 791 fn blob(&mut self) -> Result<&'a [u8], Error> { 792 let length = usize::try_from(self.u32()?).map_err(|_| Error::CorruptOutboxRecord)?; 793 self.take(length) 794 } 795 796 fn string_blob(&mut self) -> Result<&'a str, Error> { 797 core::str::from_utf8(self.blob()?).map_err(|_| Error::CorruptOutboxRecord) 798 } 799 800 fn array<const N: usize>(&mut self) -> Result<[u8; N], Error> { 801 self.take(N)? 802 .try_into() 803 .map_err(|_| Error::CorruptOutboxRecord) 804 } 805 806 fn take(&mut self, length: usize) -> Result<&'a [u8], Error> { 807 let end = self 808 .offset 809 .checked_add(length) 810 .ok_or(Error::CorruptOutboxRecord)?; 811 let value = self 812 .bytes 813 .get(self.offset..end) 814 .ok_or(Error::CorruptOutboxRecord)?; 815 self.offset = end; 816 Ok(value) 817 } 818 819 fn finish(self) -> Result<(), Error> { 820 if self.offset == self.bytes.len() { 821 Ok(()) 822 } else { 823 Err(Error::CorruptOutboxRecord) 824 } 825 } 826 } 827 828 const fn stage_name(value: OutboxStage) -> &'static str { 829 match value { 830 OutboxStage::Pending => "pending", 831 OutboxStage::Leased => "leased", 832 OutboxStage::Retryable => "retryable", 833 OutboxStage::Satisfied => "satisfied", 834 OutboxStage::Exhausted => "exhausted", 835 } 836 } 837 838 const fn stage(value: &str) -> Result<OutboxStage, Error> { 839 match value.as_bytes() { 840 b"pending" => Ok(OutboxStage::Pending), 841 b"leased" => Ok(OutboxStage::Leased), 842 b"retryable" => Ok(OutboxStage::Retryable), 843 b"satisfied" => Ok(OutboxStage::Satisfied), 844 b"exhausted" => Ok(OutboxStage::Exhausted), 845 _ => Err(Error::CorruptOutboxRecord), 846 } 847 } 848 849 const fn satisfaction_name(value: SatisfactionResult) -> &'static str { 850 match value { 851 SatisfactionResult::Pending => "pending", 852 SatisfactionResult::Satisfied => "satisfied", 853 SatisfactionResult::Exhausted => "exhausted", 854 } 855 } 856 857 const fn satisfaction(value: &str) -> Result<SatisfactionResult, Error> { 858 match value.as_bytes() { 859 b"pending" => Ok(SatisfactionResult::Pending), 860 b"satisfied" => Ok(SatisfactionResult::Satisfied), 861 b"exhausted" => Ok(SatisfactionResult::Exhausted), 862 _ => Err(Error::CorruptOutboxRecord), 863 } 864 } 865 866 const fn retryability_name(value: Retryability) -> &'static str { 867 match value { 868 Retryability::Retryable => "retryable", 869 Retryability::Terminal => "terminal", 870 Retryability::NotApplicable => "not_applicable", 871 } 872 } 873 874 const fn retryability_byte(value: Retryability) -> u8 { 875 match value { 876 Retryability::NotApplicable => 0, 877 Retryability::Retryable => 1, 878 Retryability::Terminal => 2, 879 } 880 } 881 882 const fn retryability(value: u8) -> Result<Retryability, Error> { 883 match value { 884 0 => Ok(Retryability::NotApplicable), 885 1 => Ok(Retryability::Retryable), 886 2 => Ok(Retryability::Terminal), 887 _ => Err(Error::CorruptOutboxRecord), 888 } 889 } 890 891 fn count(row: &sqlx::sqlite::SqliteRow, column: &str) -> Result<u64, Error> { 892 u64_from_i64(row.try_get::<i64, _>(column).map_err(map_corrupt)?) 893 } 894 895 fn array<const N: usize>(bytes: Vec<u8>) -> Result<[u8; N], Error> { 896 bytes.try_into().map_err(|_| Error::CorruptOutboxRecord) 897 } 898 899 fn i64_from_u64(value: u64) -> Result<i64, Error> { 900 i64::try_from(value).map_err(|_| Error::CorruptOutboxRecord) 901 } 902 903 fn u64_from_i64(value: i64) -> Result<u64, Error> { 904 u64::try_from(value).map_err(|_| Error::CorruptOutboxRecord) 905 } 906 907 fn map_corrupt(_: sqlx::Error) -> Error { 908 Error::CorruptOutboxRecord 909 } 910 911 #[cfg(test)] 912 #[cfg_attr(coverage_nightly, coverage(off))] 913 mod tests { 914 use super::*; 915 use crate::migration::runtime::{MIGRATIONS, migration_sql}; 916 use radroots_event::{SignedEvent, wire::Nip01EventWire}; 917 use radroots_storage::{ 918 Journal, Outbox, 919 event::SourceGeneration, 920 journal::{ 921 IdempotencyDigest, IdempotencyKey, OperationId, OperationInstanceId, PrepareOperation, 922 }, 923 outbox::{DeliveryReceipt, DeliveryTargetReceipt}, 924 status::EventStoreMode, 925 }; 926 use sqlx::sqlite::SqlitePoolOptions; 927 928 async fn store(mode: EventStoreMode) -> SqliteStorage { 929 let pool = SqlitePoolOptions::new() 930 .max_connections(1) 931 .connect("sqlite::memory:") 932 .await 933 .expect("memory SQLite"); 934 sqlx::query("PRAGMA foreign_keys = ON") 935 .execute(&pool) 936 .await 937 .expect("foreign keys"); 938 for migration in MIGRATIONS { 939 sqlx::raw_sql(migration_sql(migration.version()).expect("registered SQL")) 940 .execute(&pool) 941 .await 942 .expect("runtime migration"); 943 } 944 SqliteStorage::new( 945 pool, 946 SourceGeneration::new([31; 32]).expect("generation"), 947 mode, 948 ) 949 } 950 951 fn instance(byte: u8) -> OperationInstanceId { 952 OperationInstanceId::new([byte; 16]).expect("operation instance") 953 } 954 955 async fn prepare(store: &SqliteStorage, byte: u8) { 956 Journal::prepare( 957 store, 958 PrepareOperation::new( 959 instance(byte), 960 OperationId::SyncPush, 961 IdempotencyKey::parse(format!("outbox-operation-{byte:02x}")) 962 .expect("idempotency key"), 963 IdempotencyDigest::new([byte; 32]), 964 10, 965 ) 966 .expect("prepare operation"), 967 ) 968 .await 969 .expect("prepared journal operation"); 970 } 971 972 fn signed_event() -> SignedEvent { 973 let mut wire = Nip01EventWire { 974 id: "0".repeat(64), 975 pubkey: "585591529da0bab31b3b1b1f986611cf5f435dca84f978c89ee8a40cca7103df".to_owned(), 976 created_at: 1_800_000_100, 977 kind: 0, 978 tags: vec![], 979 content: "{\"display_name\":\"Moss Street Farm\",\"bot\":false}".to_owned(), 980 sig: "42".repeat(64), 981 extra: Default::default(), 982 }; 983 wire.id = wire 984 .computed_event_id() 985 .expect("canonical event id") 986 .to_hex(); 987 let raw_json = serde_json::json!({ 988 "id": &wire.id, 989 "pubkey": &wire.pubkey, 990 "created_at": wire.created_at, 991 "kind": wire.kind, 992 "tags": &wire.tags, 993 "content": &wire.content, 994 "sig": &wire.sig, 995 }) 996 .to_string(); 997 SignedEvent::from_wire_verified_id(wire, raw_json).expect("signed event") 998 } 999 1000 fn request() -> DeliveryRequest { 1001 DeliveryRequest::new( 1002 "sqlite-outbox-request", 1003 DeliveryPayload::new(signed_event()), 1004 TargetSet::new(vec![ 1005 Target::nostr_relay("wss://one.example").expect("first target"), 1006 Target::nostr_relay("wss://two.example").expect("second target"), 1007 ]) 1008 .expect("target set"), 1009 SatisfactionPolicy::new(SatisfactionClass::Accepted, TargetPolicy::all()), 1010 10_000, 1011 ) 1012 .expect("delivery request") 1013 } 1014 1015 fn enqueue(item: u8, operation: u8, digest: u8) -> EnqueueOutboxItem { 1016 EnqueueOutboxItem::new( 1017 OutboxItemId::new([item; 16]).expect("item id"), 1018 instance(operation), 1019 DeliveryPlanDigest::new([digest; 32]), 1020 request(), 1021 20, 1022 ) 1023 .expect("enqueue") 1024 } 1025 1026 fn claim_request(now: u64, expires: u64, seed: u8) -> ClaimOutboxItems { 1027 ClaimOutboxItems::new( 1028 LeaseOwner::parse("sqlite-worker").expect("owner"), 1029 LeaseId::new([seed; 16]).expect("lease seed"), 1030 now, 1031 expires, 1032 10, 1033 ) 1034 .expect("claim request") 1035 } 1036 1037 fn receipt(request: &DeliveryRequest, outcomes: [DeliveryOutcome; 2]) -> DeliveryReceipt { 1038 DeliveryReceipt::for_request( 1039 request, 1040 request 1041 .target_set() 1042 .targets() 1043 .iter() 1044 .cloned() 1045 .zip(outcomes) 1046 .map(|(target, outcome)| DeliveryTargetReceipt::attempted(target, outcome)) 1047 .collect(), 1048 ) 1049 .expect("delivery receipt") 1050 } 1051 1052 #[tokio::test] 1053 async fn enqueue_replays_exact_plans_and_rejects_identity_conflicts() { 1054 let store = store(EventStoreMode::ReadWrite).await; 1055 prepare(&store, 1).await; 1056 let item = enqueue(1, 1, 2); 1057 let created = store.enqueue(item.clone()).await.expect("created"); 1058 assert_eq!(created.disposition(), EnqueueDisposition::Created); 1059 let replay = store.enqueue(item).await.expect("replay"); 1060 assert_eq!(replay.disposition(), EnqueueDisposition::Replay); 1061 assert_eq!(replay.record(), created.record()); 1062 assert_eq!( 1063 store.enqueue(enqueue(1, 1, 3)).await, 1064 Err(Error::OutboxPlanConflict) 1065 ); 1066 assert_eq!( 1067 store.enqueue(enqueue(2, 1, 2)).await, 1068 Err(Error::OutboxPlanConflict) 1069 ); 1070 assert_eq!( 1071 store 1072 .item(OutboxItemId::new([1; 16]).expect("item id")) 1073 .await 1074 .expect("lookup") 1075 .expect("record"), 1076 *created.record() 1077 ); 1078 } 1079 1080 #[tokio::test] 1081 async fn claims_expire_release_and_honor_retry_deferral() { 1082 let store = store(EventStoreMode::ReadWrite).await; 1083 prepare(&store, 2).await; 1084 store.enqueue(enqueue(2, 2, 2)).await.expect("enqueue"); 1085 let first = store 1086 .claim(claim_request(100, 200, 3)) 1087 .await 1088 .expect("claim") 1089 .pop() 1090 .expect("claimed"); 1091 assert!( 1092 store 1093 .claim(claim_request(150, 250, 4)) 1094 .await 1095 .expect("concurrent claim") 1096 .is_empty() 1097 ); 1098 let reclaimed = store 1099 .claim(claim_request(200, 300, 5)) 1100 .await 1101 .expect("expired claim") 1102 .pop() 1103 .expect("reclaimed"); 1104 assert_ne!(first.lease().id(), reclaimed.lease().id()); 1105 let released = store 1106 .release( 1107 reclaimed.record().item_id(), 1108 reclaimed.lease().id(), 1109 reclaimed.record().revision(), 1110 210, 1111 Some(250), 1112 ) 1113 .await 1114 .expect("release"); 1115 assert_eq!(released.stage(), OutboxStage::Pending); 1116 assert!( 1117 store 1118 .claim(claim_request(249, 300, 6)) 1119 .await 1120 .expect("deferred claim") 1121 .is_empty() 1122 ); 1123 assert_eq!( 1124 store 1125 .claim(claim_request(250, 350, 7)) 1126 .await 1127 .expect("ready claim") 1128 .len(), 1129 1 1130 ); 1131 } 1132 1133 #[tokio::test] 1134 async fn partial_evidence_survives_recovery_and_advances_to_satisfaction() { 1135 let store = store(EventStoreMode::ReadWrite).await; 1136 prepare(&store, 3).await; 1137 store.enqueue(enqueue(3, 3, 3)).await.expect("enqueue"); 1138 let first = store 1139 .claim(claim_request(100, 200, 4)) 1140 .await 1141 .expect("claim") 1142 .pop() 1143 .expect("claimed"); 1144 let partial = store 1145 .record_attempt( 1146 DeliveryAttemptEvidence::new( 1147 first.record().item_id(), 1148 first.lease().id(), 1149 first.record().revision(), 1150 DeliveryAttempt::FIRST, 1151 receipt( 1152 first.record().request(), 1153 [DeliveryOutcome::accepted(), DeliveryOutcome::unavailable()], 1154 ), 1155 150, 1156 ) 1157 .expect("partial evidence"), 1158 ) 1159 .await 1160 .expect("partial attempt"); 1161 assert_eq!(partial.stage(), OutboxStage::Retryable); 1162 assert_eq!(partial.evidence().len(), 2); 1163 1164 let recovered = SqliteStorage::new( 1165 store.pool().clone(), 1166 SourceGeneration::new([31; 32]).expect("generation"), 1167 EventStoreMode::ReadWrite, 1168 ); 1169 let second = recovered 1170 .claim(claim_request(250, 350, 5)) 1171 .await 1172 .expect("recovery claim") 1173 .pop() 1174 .expect("reclaimed"); 1175 let satisfied = recovered 1176 .record_attempt( 1177 DeliveryAttemptEvidence::new( 1178 second.record().item_id(), 1179 second.lease().id(), 1180 second.record().revision(), 1181 DeliveryAttempt::new(2).expect("second attempt"), 1182 receipt( 1183 second.record().request(), 1184 [DeliveryOutcome::accepted(), DeliveryOutcome::delivered()], 1185 ), 1186 300, 1187 ) 1188 .expect("success evidence"), 1189 ) 1190 .await 1191 .expect("successful attempt"); 1192 assert_eq!(satisfied.stage(), OutboxStage::Satisfied); 1193 assert_eq!(satisfied.evidence().len(), 4); 1194 assert_eq!(recovered.status().await.expect("status").satisfied, 1); 1195 1196 sqlx::query( 1197 "UPDATE radroots_runtime_delivery_evidence SET outcome = X'FF' 1198 WHERE item_id = ? AND attempt = 2", 1199 ) 1200 .bind(satisfied.item_id().as_bytes().as_slice()) 1201 .execute(recovered.pool()) 1202 .await 1203 .expect("forge corrupt evidence"); 1204 assert_eq!( 1205 recovered.item(satisfied.item_id()).await, 1206 Err(Error::CorruptOutboxRecord) 1207 ); 1208 let read_only = SqliteStorage::new( 1209 recovered.pool().clone(), 1210 SourceGeneration::new([31; 32]).expect("generation"), 1211 EventStoreMode::ReadOnly, 1212 ); 1213 assert_eq!( 1214 read_only.claim(claim_request(400, 500, 6)).await, 1215 Err(Error::BackendUnavailable) 1216 ); 1217 } 1218 1219 #[test] 1220 fn versioned_codecs_round_trip_every_policy_target_and_outcome_shape() { 1221 let targets = TargetSet::new(vec![ 1222 Target::local_with_metadata( 1223 "local://queue", 1224 Some(TargetScope::parse("farm.one").expect("scope")), 1225 Some(TargetLabel::parse("Farm queue").expect("label")), 1226 ) 1227 .expect("local target"), 1228 Target::nostr_relay("wss://relay.example").expect("relay target"), 1229 ]) 1230 .expect("target set"); 1231 let required = targets.targets()[0].fingerprint().clone(); 1232 for policy in [ 1233 TargetPolicy::any(), 1234 TargetPolicy::all(), 1235 TargetPolicy::quorum(1).expect("quorum"), 1236 TargetPolicy::required(vec![required]).expect("required"), 1237 ] { 1238 let request = DeliveryRequest::new( 1239 "codec-request", 1240 DeliveryPayload::new(signed_event()), 1241 targets.clone(), 1242 SatisfactionPolicy::new(SatisfactionClass::Delivered, policy), 1243 20_000, 1244 ) 1245 .expect("request"); 1246 assert_eq!( 1247 decode_request(&encode_request(&request).expect("encode request")) 1248 .expect("decode request"), 1249 request 1250 ); 1251 } 1252 1253 let outcomes = [ 1254 DeliveryOutcome::accepted(), 1255 DeliveryOutcome::delivered(), 1256 DeliveryOutcome::rejected(), 1257 DeliveryOutcome::unavailable(), 1258 DeliveryOutcome::failed(Retryability::Retryable) 1259 .expect("retryable failure") 1260 .with_detail("relay_timeout", "relay request timed out") 1261 .expect("failure detail"), 1262 DeliveryOutcome::failed(Retryability::Terminal).expect("terminal failure"), 1263 ]; 1264 for outcome in outcomes { 1265 assert_eq!( 1266 decode_outcome(&encode_outcome(&outcome).expect("encode outcome")) 1267 .expect("decode outcome"), 1268 outcome 1269 ); 1270 } 1271 for target in targets.targets() { 1272 let encoded = encode_target(target).expect("encode target"); 1273 assert_eq!(decode_target(&encoded).expect("decode target"), *target); 1274 assert_eq!(decode_target(&[0]), Err(Error::CorruptOutboxRecord)); 1275 let mut trailing = encoded; 1276 trailing.push(0); 1277 assert_eq!(decode_target(&trailing), Err(Error::CorruptOutboxRecord)); 1278 } 1279 1280 let encoded = encode_request(&request()).expect("encode request"); 1281 for end in 0..encoded.len() { 1282 let _ = decode_request(&encoded[..end]); 1283 } 1284 let mut invalid_version = encoded.clone(); 1285 invalid_version[0] = 0; 1286 assert_eq!( 1287 decode_request(&invalid_version), 1288 Err(Error::CorruptOutboxRecord) 1289 ); 1290 let request_id_length = usize::from(u16::from_be_bytes([encoded[1], encoded[2]])); 1291 let event_length_offset = 3 + request_id_length; 1292 let event_length = usize::try_from(u32::from_be_bytes( 1293 encoded[event_length_offset..event_length_offset + 4] 1294 .try_into() 1295 .expect("event length"), 1296 )) 1297 .expect("event length fits"); 1298 let target_count_offset = event_length_offset + 4 + event_length; 1299 let mut no_targets = encoded.clone(); 1300 no_targets[target_count_offset..target_count_offset + 2] 1301 .copy_from_slice(&0_u16.to_be_bytes()); 1302 assert_eq!(decode_request(&no_targets), Err(Error::CorruptOutboxRecord)); 1303 let mut too_many_targets = encoded; 1304 too_many_targets[target_count_offset..target_count_offset + 2] 1305 .copy_from_slice(&u16::MAX.to_be_bytes()); 1306 assert_eq!( 1307 decode_request(&too_many_targets), 1308 Err(Error::CorruptOutboxRecord) 1309 ); 1310 1311 for kind in 0..=5 { 1312 for retryability in 0..=3 { 1313 for detail in 0..=2 { 1314 let _ = decode_outcome(&[1, kind, retryability, detail]); 1315 } 1316 } 1317 } 1318 assert_eq!(decode_outcome(&[0]), Err(Error::CorruptOutboxRecord)); 1319 } 1320 }