blossom.rs (89725B)
1 use std::{net::SocketAddr, time::Duration}; 2 3 use radroots_blossom::{BlobDescriptor, BlobUrl, MediaType, Sha256, descriptor::ByteCommitment}; 4 use reqwest::{ 5 StatusCode, 6 header::{ 7 ACCEPT, ACCEPT_ENCODING, AUTHORIZATION, CONTENT_ENCODING, CONTENT_LENGTH, CONTENT_TYPE, 8 LOCATION, 9 }, 10 }; 11 12 use crate::transport::{ 13 BlossomCancellation, BlossomConfig, BlossomEndpoint, BlossomError, BlossomErrorKind, 14 BlossomImageDimensions, BlossomInboundReceipt, BlossomInboundRequest, BlossomPhase, 15 BlossomUploadReceipt, BlossomUploadRequest, BlossomUploadTransaction, 16 }; 17 18 const MAX_RESOLVED_ADDRESSES: usize = 32; 19 const MAX_ERROR_RESPONSE_BYTES: usize = 4_096; 20 const MAX_SERVER_ERROR_CODE_BYTES: usize = 64; 21 const X_SHA_256: &str = "x-sha-256"; 22 const BLOSSOM_PROBE_HASH: &str = "0000000000000000000000000000000000000000000000000000000000000000"; 23 24 pub(crate) struct BlossomProbeObservation { 25 pub(crate) http_status: u16, 26 } 27 28 pub(crate) struct BlossomProbeFailure { 29 pub(crate) error: BlossomError, 30 pub(crate) dns_policy_validated: bool, 31 } 32 33 /// Performs one BUD-01-shaped GET for an impossible sentinel digest. 34 /// 35 /// The request carries no authorization and cannot upload, delete, or mutate a 36 /// server. Any terminal HTTP response proves only DNS-policy, transport, and 37 /// HTTP reachability for the exact configured origin. 38 #[cfg_attr(coverage_nightly, coverage(off))] 39 pub(crate) async fn probe( 40 config: BlossomConfig, 41 endpoint: BlossomEndpoint, 42 cancellation: BlossomCancellation, 43 ) -> Result<BlossomProbeObservation, BlossomProbeFailure> { 44 let mut url = BlobUrl::parse(format!("{}/{BLOSSOM_PROBE_HASH}", endpoint.origin()).as_str()) 45 .map_err(|_| BlossomProbeFailure { 46 error: failure( 47 BlossomErrorKind::InvalidEndpoint, 48 BlossomPhase::Probe, 49 false, 50 false, 51 0, 52 ), 53 dns_policy_validated: false, 54 })?; 55 let mut dns_policy_validated = false; 56 for redirects in 0..=config.max_redirects() { 57 let current = 58 config 59 .profile() 60 .endpoint_for_blob(&url) 61 .ok_or_else(|| BlossomProbeFailure { 62 error: failure( 63 BlossomErrorKind::UnsafeRedirect, 64 BlossomPhase::Probe, 65 false, 66 false, 67 1, 68 ), 69 dns_policy_validated, 70 })?; 71 let addresses = resolve( 72 current, 73 config.connect_timeout(), 74 &cancellation, 75 BlossomPhase::Probe, 76 1, 77 false, 78 ) 79 .await 80 .map_err(|error| BlossomProbeFailure { 81 error, 82 dns_policy_validated, 83 })?; 84 dns_policy_validated = true; 85 let client = hardened_client_with_addresses( 86 &config, 87 current, 88 addresses.as_slice(), 89 BlossomPhase::Probe, 90 1, 91 false, 92 ) 93 .map_err(|error| BlossomProbeFailure { 94 error, 95 dns_policy_validated, 96 })?; 97 let pending = client 98 .get(url.as_str()) 99 .header(ACCEPT, "*/*") 100 .header(ACCEPT_ENCODING, "identity") 101 .send(); 102 let response = tokio::select! { 103 biased; 104 _ = cancellation.cancelled() => { 105 return Err(BlossomProbeFailure { 106 error: failure( 107 BlossomErrorKind::Cancelled, 108 BlossomPhase::Probe, 109 true, 110 false, 111 1, 112 ), 113 dns_policy_validated, 114 }); 115 } 116 response = pending => response.map_err(|error| BlossomProbeFailure { 117 error: request_error(error, BlossomPhase::Probe, false, 1), 118 dns_policy_validated, 119 })?, 120 }; 121 if response.status().is_redirection() { 122 if redirects == config.max_redirects() { 123 return Err(BlossomProbeFailure { 124 error: failure( 125 BlossomErrorKind::RedirectLimit, 126 BlossomPhase::Probe, 127 false, 128 false, 129 1, 130 ), 131 dns_policy_validated, 132 }); 133 } 134 let location = response 135 .headers() 136 .get(LOCATION) 137 .and_then(|value| value.to_str().ok()) 138 .ok_or_else(|| BlossomProbeFailure { 139 error: failure( 140 BlossomErrorKind::UnsafeRedirect, 141 BlossomPhase::Probe, 142 false, 143 false, 144 1, 145 ), 146 dns_policy_validated, 147 })?; 148 let base = reqwest::Url::parse(url.as_str()).map_err(|_| BlossomProbeFailure { 149 error: failure( 150 BlossomErrorKind::UnsafeRedirect, 151 BlossomPhase::Probe, 152 false, 153 false, 154 1, 155 ), 156 dns_policy_validated, 157 })?; 158 let next = base 159 .join(location) 160 .ok() 161 .and_then(|value| BlobUrl::parse(value.as_str()).ok()) 162 .filter(|value| { 163 value.hash_path().hash() == url.hash_path().hash() 164 && config.profile().endpoint_for_blob(value).is_some() 165 }) 166 .ok_or_else(|| BlossomProbeFailure { 167 error: failure( 168 BlossomErrorKind::UnsafeRedirect, 169 BlossomPhase::Probe, 170 false, 171 false, 172 1, 173 ), 174 dns_policy_validated, 175 })?; 176 url = next; 177 continue; 178 } 179 let status = response.status().as_u16(); 180 read_bounded( 181 response, 182 config.max_descriptor_bytes(), 183 &cancellation, 184 BlossomPhase::Probe, 185 false, 186 1, 187 ) 188 .await 189 .map_err(|error| BlossomProbeFailure { 190 error, 191 dns_policy_validated, 192 })?; 193 return Ok(BlossomProbeObservation { 194 http_status: status, 195 }); 196 } 197 Err(BlossomProbeFailure { 198 error: failure( 199 BlossomErrorKind::RedirectLimit, 200 BlossomPhase::Probe, 201 false, 202 false, 203 1, 204 ), 205 dns_policy_validated, 206 }) 207 } 208 209 pub(crate) async fn upload( 210 transaction: BlossomUploadTransaction, 211 authorization: crate::signing::AuthorizationHeader, 212 cancellation: BlossomCancellation, 213 ) -> Result<BlossomUploadReceipt, BlossomError> { 214 upload_bound_with_authorization(transaction, authorization.as_str(), cancellation).await 215 } 216 217 async fn upload_bound_with_authorization( 218 transaction: BlossomUploadTransaction, 219 authorization: &str, 220 cancellation: BlossomCancellation, 221 ) -> Result<BlossomUploadReceipt, BlossomError> { 222 let config = transaction.config().clone(); 223 let endpoint = transaction.endpoint().clone(); 224 let expected_url = transaction.expected_url().clone(); 225 let request = transaction.into_request(); 226 if request.byte_size() > config.max_blob_bytes() { 227 return Err(failure( 228 BlossomErrorKind::ResponseTooLarge, 229 BlossomPhase::Verification, 230 false, 231 false, 232 0, 233 )); 234 } 235 let mut upload_attempts = 0_u8; 236 let descriptor = loop { 237 ensure_not_cancelled(&cancellation, BlossomPhase::Upload, upload_attempts, false)?; 238 upload_attempts = upload_attempts.saturating_add(1); 239 match upload_once( 240 &config, 241 &endpoint, 242 &request, 243 authorization, 244 &cancellation, 245 upload_attempts, 246 ) 247 .await 248 { 249 Ok(descriptor) => break descriptor, 250 Err(error) if error.retryable() && upload_attempts < config.max_attempts() => { 251 retry_delay( 252 &config, 253 upload_attempts, 254 &cancellation, 255 BlossomPhase::Upload, 256 true, 257 ) 258 .await?; 259 } 260 Err(error) => return Err(error), 261 } 262 }; 263 264 let verified_upload = verify_descriptor(&request, &expected_url, descriptor, upload_attempts)?; 265 let mut retrieval_attempts = 0_u8; 266 let retrieved = loop { 267 ensure_not_cancelled( 268 &cancellation, 269 BlossomPhase::Retrieval, 270 upload_attempts.saturating_add(retrieval_attempts), 271 true, 272 )?; 273 retrieval_attempts = retrieval_attempts.saturating_add(1); 274 match retrieve_once( 275 &config, 276 verified_upload.url().as_blob_url().clone(), 277 request.sha256(), 278 Some(request.media_type()), 279 Some(request.byte_size()), 280 Some(request.dimensions()), 281 &cancellation, 282 upload_attempts.saturating_add(retrieval_attempts), 283 true, 284 ) 285 .await 286 { 287 Ok(retrieved) => break retrieved.bytes, 288 Err(error) if error.retryable() && retrieval_attempts < config.max_attempts() => { 289 retry_delay( 290 &config, 291 retrieval_attempts, 292 &cancellation, 293 BlossomPhase::Retrieval, 294 true, 295 ) 296 .await?; 297 } 298 Err(error) => return Err(error), 299 } 300 }; 301 302 if retrieved.as_slice() != request.bytes() { 303 return Err(failure( 304 BlossomErrorKind::RetrievedBytesMismatch, 305 BlossomPhase::Verification, 306 false, 307 true, 308 upload_attempts.saturating_add(retrieval_attempts), 309 )); 310 } 311 verify_image( 312 retrieved.as_slice(), 313 request.media_type(), 314 request.dimensions(), 315 ) 316 .map_err(|error| with_operation(error, true, upload_attempts + retrieval_attempts))?; 317 let verified_retrieval = verified_upload 318 .into_descriptor() 319 .approve_reference() 320 .and_then(|approved| approved.verify_bytes(retrieved.as_slice(), request.media_type())) 321 .map_err(|_| { 322 failure( 323 BlossomErrorKind::RetrievedBytesMismatch, 324 BlossomPhase::Verification, 325 false, 326 true, 327 upload_attempts.saturating_add(retrieval_attempts), 328 ) 329 })?; 330 331 Ok(BlossomUploadReceipt::new( 332 verified_retrieval, 333 request.dimensions(), 334 upload_attempts.saturating_add(retrieval_attempts), 335 request.verified_at_unix_ms(), 336 )) 337 } 338 339 pub(crate) async fn complete_native_upload( 340 transaction: BlossomUploadTransaction, 341 status_code: u16, 342 response_media_type: Option<&str>, 343 response_content_encoding: Option<&str>, 344 response_body: &[u8], 345 cancellation: BlossomCancellation, 346 ) -> Result<BlossomUploadReceipt, BlossomError> { 347 let config = transaction.config().clone(); 348 let expected_url = transaction.expected_url().clone(); 349 let request = transaction.into_request(); 350 if !matches!(status_code, 200 | 201) { 351 let status = StatusCode::from_u16(status_code).unwrap_or(StatusCode::INTERNAL_SERVER_ERROR); 352 return Err(http_status_error(status, BlossomPhase::Upload, true, 1)); 353 } 354 if response_media_type != Some("application/json") 355 || response_content_encoding.is_some_and(|value| value != "identity") 356 || response_body.is_empty() 357 || response_body.len() > config.max_descriptor_bytes() 358 { 359 return Err(failure( 360 BlossomErrorKind::InvalidDescriptor, 361 BlossomPhase::Upload, 362 false, 363 true, 364 1, 365 )); 366 } 367 let descriptor = serde_json::from_slice::<BlobDescriptor>(response_body).map_err(|_| { 368 failure( 369 BlossomErrorKind::InvalidDescriptor, 370 BlossomPhase::Upload, 371 false, 372 true, 373 1, 374 ) 375 })?; 376 let verified_upload = verify_descriptor(&request, &expected_url, descriptor, 1)?; 377 let inbound = BlossomInboundRequest::new( 378 verified_upload.url().as_blob_url().clone(), 379 Some(request.media_type().clone()), 380 Some(request.byte_size()), 381 Some(request.dimensions()), 382 ) 383 .map_err(|error| with_operation(error, true, 1))?; 384 let retrieved = retrieve(config, inbound, cancellation) 385 .await 386 .map_err(|error| { 387 let attempts = 1_u8.saturating_add(error.attempts()); 388 with_operation(error, true, attempts) 389 })?; 390 if retrieved.bytes() != request.bytes() { 391 return Err(failure( 392 BlossomErrorKind::RetrievedBytesMismatch, 393 BlossomPhase::Verification, 394 false, 395 true, 396 1_u8.saturating_add(retrieved.attempts()), 397 )); 398 } 399 let descriptor = verified_upload 400 .into_descriptor() 401 .approve_reference() 402 .and_then(|approved| approved.verify_bytes(retrieved.bytes(), request.media_type())) 403 .map_err(|_| { 404 failure( 405 BlossomErrorKind::RetrievedBytesMismatch, 406 BlossomPhase::Verification, 407 false, 408 true, 409 1_u8.saturating_add(retrieved.attempts()), 410 ) 411 })?; 412 Ok(BlossomUploadReceipt::new( 413 descriptor, 414 request.dimensions(), 415 1_u8.saturating_add(retrieved.attempts()), 416 request.verified_at_unix_ms(), 417 )) 418 } 419 420 pub(crate) async fn retrieve( 421 config: BlossomConfig, 422 request: BlossomInboundRequest, 423 cancellation: BlossomCancellation, 424 ) -> Result<BlossomInboundReceipt, BlossomError> { 425 if request 426 .expected_byte_size() 427 .is_some_and(|size| size > config.max_blob_bytes()) 428 { 429 return Err(failure( 430 BlossomErrorKind::ResponseTooLarge, 431 BlossomPhase::Verification, 432 false, 433 false, 434 0, 435 )); 436 } 437 let fingerprint = config.fingerprint(); 438 let mut attempts = 0_u8; 439 let retrieved = loop { 440 ensure_not_cancelled(&cancellation, BlossomPhase::Retrieval, attempts, false)?; 441 attempts = attempts.saturating_add(1); 442 match retrieve_once( 443 &config, 444 request.url().clone(), 445 request.url().hash_path().hash(), 446 request.expected_media_type(), 447 request.expected_byte_size(), 448 request.expected_dimensions(), 449 &cancellation, 450 attempts, 451 false, 452 ) 453 .await 454 { 455 Ok(retrieved) => break retrieved, 456 Err(error) if error.retryable() && attempts < config.max_attempts() => { 457 retry_delay( 458 &config, 459 attempts, 460 &cancellation, 461 BlossomPhase::Retrieval, 462 false, 463 ) 464 .await?; 465 } 466 Err(error) => return Err(error), 467 } 468 }; 469 let commitment = ByteCommitment::from_bytes(retrieved.bytes.as_slice(), retrieved.media_type); 470 Ok(BlossomInboundReceipt::new( 471 retrieved.final_url, 472 commitment, 473 retrieved.dimensions, 474 std::sync::Arc::from(retrieved.bytes), 475 fingerprint, 476 attempts, 477 crate::transport::blossom_now_unix_ms(), 478 )) 479 } 480 481 #[cfg(test)] 482 async fn upload_with_authorization( 483 config: BlossomConfig, 484 request: BlossomUploadRequest, 485 authorization: &str, 486 cancellation: BlossomCancellation, 487 ) -> Result<BlossomUploadReceipt, BlossomError> { 488 let slot = crate::transport::BlossomSlot::new(); 489 slot.configure(config)?; 490 let transaction = slot.prepare_upload(request)?; 491 upload_bound_with_authorization(transaction, authorization, cancellation).await 492 } 493 494 // Direct DNS/socket/HTTP behavior is verified by the local real-I/O suite; 495 // deterministic coverage owns the surrounding retry, validation, and durable 496 // state policy. 497 #[cfg_attr(coverage_nightly, coverage(off))] 498 async fn upload_once( 499 config: &BlossomConfig, 500 endpoint: &BlossomEndpoint, 501 request: &BlossomUploadRequest, 502 authorization: &str, 503 cancellation: &BlossomCancellation, 504 attempt: u8, 505 ) -> Result<BlobDescriptor, BlossomError> { 506 let client = hardened_client( 507 config, 508 endpoint, 509 cancellation, 510 BlossomPhase::Upload, 511 attempt, 512 false, 513 ) 514 .await?; 515 let pending = client 516 .put(endpoint.upload_url()) 517 .header(AUTHORIZATION, authorization) 518 .header(X_SHA_256, request.sha256().to_string()) 519 .header(CONTENT_TYPE, request.media_type().as_str()) 520 .header(ACCEPT, "application/json") 521 .header(ACCEPT_ENCODING, "identity") 522 .body(request.bytes().to_vec()) 523 .send(); 524 let response = tokio::select! { 525 biased; 526 _ = cancellation.cancelled() => { 527 return Err(failure( 528 BlossomErrorKind::Cancelled, 529 BlossomPhase::Upload, 530 true, 531 true, 532 attempt, 533 )); 534 } 535 response = pending => response.map_err(|error| request_error(error, BlossomPhase::Upload, true, attempt))?, 536 }; 537 538 if response.status().is_redirection() { 539 return Err(failure( 540 BlossomErrorKind::UnsafeRedirect, 541 BlossomPhase::Upload, 542 false, 543 true, 544 attempt, 545 )); 546 } 547 if !matches!(response.status(), StatusCode::OK | StatusCode::CREATED) { 548 return Err(http_status_response_error( 549 response, 550 BlossomPhase::Upload, 551 true, 552 attempt, 553 cancellation, 554 ) 555 .await); 556 } 557 { 558 let content_type = response 559 .headers() 560 .get(CONTENT_TYPE) 561 .and_then(|value| value.to_str().ok()) 562 .ok_or_else(|| { 563 failure( 564 BlossomErrorKind::InvalidDescriptor, 565 BlossomPhase::Descriptor, 566 false, 567 true, 568 attempt, 569 ) 570 })?; 571 if content_type 572 .split(';') 573 .next() 574 .is_none_or(|value| !value.trim().eq_ignore_ascii_case("application/json")) 575 { 576 return Err(failure( 577 BlossomErrorKind::InvalidDescriptor, 578 BlossomPhase::Descriptor, 579 false, 580 true, 581 attempt, 582 )); 583 } 584 } 585 let bytes = read_bounded( 586 response, 587 config.max_descriptor_bytes(), 588 cancellation, 589 BlossomPhase::Descriptor, 590 true, 591 attempt, 592 ) 593 .await?; 594 serde_json::from_slice(bytes.as_slice()).map_err(|_| { 595 failure( 596 BlossomErrorKind::InvalidDescriptor, 597 BlossomPhase::Descriptor, 598 false, 599 true, 600 attempt, 601 ) 602 }) 603 } 604 605 fn verify_descriptor( 606 request: &BlossomUploadRequest, 607 expected_url: &BlobUrl, 608 descriptor: BlobDescriptor, 609 attempts: u8, 610 ) -> Result<radroots_blossom::ByteVerifiedDescriptor, BlossomError> { 611 // `BlobDescriptor` construction already binds `sha256` to the URL hash, 612 // so equality of the typed URL proves equality of that hash as well. 613 if descriptor.url() != expected_url 614 || descriptor.size() != request.byte_size() 615 || descriptor.media_type() != request.media_type() 616 { 617 return Err(failure( 618 BlossomErrorKind::DescriptorMismatch, 619 BlossomPhase::Descriptor, 620 false, 621 true, 622 attempts, 623 )); 624 } 625 descriptor 626 .approve_reference() 627 .and_then(|approved| approved.verify_bytes(request.bytes(), request.media_type())) 628 .map_err(|_| { 629 failure( 630 BlossomErrorKind::DescriptorMismatch, 631 BlossomPhase::Descriptor, 632 false, 633 true, 634 attempts, 635 ) 636 }) 637 } 638 639 struct RetrievedBytes { 640 final_url: BlobUrl, 641 bytes: Vec<u8>, 642 media_type: MediaType, 643 dimensions: BlossomImageDimensions, 644 } 645 646 #[cfg_attr(coverage_nightly, coverage(off))] 647 #[allow(clippy::too_many_arguments)] 648 async fn retrieve_once( 649 config: &BlossomConfig, 650 mut url: BlobUrl, 651 expected_sha256: Sha256, 652 expected_media_type: Option<&MediaType>, 653 expected_byte_size: Option<u64>, 654 expected_dimensions: Option<BlossomImageDimensions>, 655 cancellation: &BlossomCancellation, 656 attempt: u8, 657 possible_orphan: bool, 658 ) -> Result<RetrievedBytes, BlossomError> { 659 for redirects in 0..=config.max_redirects() { 660 let endpoint = config.profile().endpoint_for_blob(&url).ok_or_else(|| { 661 failure( 662 BlossomErrorKind::UnsafeRedirect, 663 BlossomPhase::Retrieval, 664 false, 665 possible_orphan, 666 attempt, 667 ) 668 })?; 669 let client = hardened_client( 670 config, 671 endpoint, 672 cancellation, 673 BlossomPhase::Retrieval, 674 attempt, 675 possible_orphan, 676 ) 677 .await?; 678 let pending = client 679 .get(url.as_str()) 680 .header( 681 ACCEPT, 682 expected_media_type.map_or("image/*", MediaType::as_str), 683 ) 684 .header(ACCEPT_ENCODING, "identity") 685 .send(); 686 let response = tokio::select! { 687 biased; 688 _ = cancellation.cancelled() => { 689 return Err(failure( 690 BlossomErrorKind::Cancelled, 691 BlossomPhase::Retrieval, 692 true, 693 possible_orphan, 694 attempt, 695 )); 696 } 697 response = pending => response.map_err(|error| request_error(error, BlossomPhase::Retrieval, possible_orphan, attempt))?, 698 }; 699 if response.status().is_redirection() { 700 if redirects == config.max_redirects() { 701 return Err(failure( 702 BlossomErrorKind::RedirectLimit, 703 BlossomPhase::Retrieval, 704 false, 705 possible_orphan, 706 attempt, 707 )); 708 } 709 let location = response 710 .headers() 711 .get(LOCATION) 712 .and_then(|value| value.to_str().ok()) 713 .ok_or_else(|| { 714 failure( 715 BlossomErrorKind::UnsafeRedirect, 716 BlossomPhase::Retrieval, 717 false, 718 possible_orphan, 719 attempt, 720 ) 721 })?; 722 let base = reqwest::Url::parse(url.as_str()).map_err(|_| { 723 failure( 724 BlossomErrorKind::UnsafeRedirect, 725 BlossomPhase::Retrieval, 726 false, 727 possible_orphan, 728 attempt, 729 ) 730 })?; 731 let next = base.join(location).map_err(|_| { 732 failure( 733 BlossomErrorKind::UnsafeRedirect, 734 BlossomPhase::Retrieval, 735 false, 736 possible_orphan, 737 attempt, 738 ) 739 })?; 740 let next = BlobUrl::parse(next.as_str()).map_err(|_| { 741 failure( 742 BlossomErrorKind::UnsafeRedirect, 743 BlossomPhase::Retrieval, 744 false, 745 possible_orphan, 746 attempt, 747 ) 748 })?; 749 if next.hash_path().hash() != expected_sha256 750 || next.clone().approve().is_err() 751 || config.profile().endpoint_for_blob(&next).is_none() 752 { 753 return Err(failure( 754 BlossomErrorKind::UnsafeRedirect, 755 BlossomPhase::Retrieval, 756 false, 757 possible_orphan, 758 attempt, 759 )); 760 } 761 url = next; 762 continue; 763 } 764 if response.status() != StatusCode::OK { 765 return Err(http_status_response_error( 766 response, 767 BlossomPhase::Retrieval, 768 possible_orphan, 769 attempt, 770 cancellation, 771 ) 772 .await); 773 } 774 let actual_media_type = response 775 .headers() 776 .get(CONTENT_TYPE) 777 .and_then(|value| value.to_str().ok()) 778 .and_then(|value| MediaType::parse(value).ok()) 779 .ok_or_else(|| { 780 failure( 781 BlossomErrorKind::MediaTypeMismatch, 782 BlossomPhase::Verification, 783 false, 784 possible_orphan, 785 attempt, 786 ) 787 })?; 788 if expected_media_type.is_some_and(|expected| expected != &actual_media_type) { 789 return Err(failure( 790 BlossomErrorKind::MediaTypeMismatch, 791 BlossomPhase::Verification, 792 false, 793 possible_orphan, 794 attempt, 795 )); 796 } 797 let canonical_extension = 798 crate::transport::canonical_image_extension(&actual_media_type) 799 .map_err(|error| with_operation(error, possible_orphan, attempt))?; 800 if url 801 .hash_path() 802 .extension() 803 .is_none_or(|value| value.as_str() != canonical_extension) 804 { 805 return Err(failure( 806 BlossomErrorKind::MediaTypeMismatch, 807 BlossomPhase::Verification, 808 false, 809 possible_orphan, 810 attempt, 811 )); 812 } 813 if response 814 .headers() 815 .get(CONTENT_ENCODING) 816 .is_some_and(|value| { 817 value 818 .to_str() 819 .map_or(true, |encoding| !encoding.eq_ignore_ascii_case("identity")) 820 }) 821 { 822 return Err(failure( 823 BlossomErrorKind::ContentEncodingDenied, 824 BlossomPhase::Retrieval, 825 false, 826 possible_orphan, 827 attempt, 828 )); 829 } 830 let content_length = match response.headers().get(CONTENT_LENGTH) { 831 Some(value) => Some( 832 value 833 .to_str() 834 .ok() 835 .and_then(|value| value.parse::<u64>().ok()) 836 .ok_or_else(|| { 837 failure( 838 BlossomErrorKind::ResponseSizeMismatch, 839 BlossomPhase::Retrieval, 840 false, 841 possible_orphan, 842 attempt, 843 ) 844 })?, 845 ), 846 None => None, 847 }; 848 if content_length.is_some_and(|size| size > config.max_blob_bytes()) { 849 return Err(failure( 850 BlossomErrorKind::ResponseTooLarge, 851 BlossomPhase::Retrieval, 852 false, 853 possible_orphan, 854 attempt, 855 )); 856 } 857 if expected_byte_size 858 .zip(content_length) 859 .is_some_and(|(expected, actual)| expected != actual) 860 { 861 return Err(failure( 862 BlossomErrorKind::ResponseSizeMismatch, 863 BlossomPhase::Verification, 864 false, 865 possible_orphan, 866 attempt, 867 )); 868 } 869 let bytes = read_bounded( 870 response, 871 usize::try_from(config.max_blob_bytes()).unwrap_or(usize::MAX), 872 cancellation, 873 BlossomPhase::Retrieval, 874 possible_orphan, 875 attempt, 876 ) 877 .await?; 878 let byte_size = u64::try_from(bytes.len()).unwrap_or(u64::MAX); 879 if expected_byte_size.is_some_and(|expected| expected != byte_size) 880 || content_length.is_some_and(|expected| expected != byte_size) 881 { 882 return Err(failure( 883 BlossomErrorKind::ResponseSizeMismatch, 884 BlossomPhase::Verification, 885 false, 886 possible_orphan, 887 attempt, 888 )); 889 } 890 if Sha256::digest(bytes.as_slice()) != expected_sha256 { 891 return Err(failure( 892 BlossomErrorKind::ResponseHashMismatch, 893 BlossomPhase::Verification, 894 false, 895 possible_orphan, 896 attempt, 897 )); 898 } 899 let dimensions = inspect_image(bytes.as_slice(), &actual_media_type) 900 .map_err(|error| with_operation(error, possible_orphan, attempt))?; 901 if expected_dimensions.is_some_and(|expected| expected != dimensions) { 902 return Err(failure( 903 BlossomErrorKind::DimensionMismatch, 904 BlossomPhase::Verification, 905 false, 906 possible_orphan, 907 attempt, 908 )); 909 } 910 return Ok(RetrievedBytes { 911 final_url: url, 912 bytes, 913 media_type: actual_media_type, 914 dimensions, 915 }); 916 } 917 Err(failure( 918 BlossomErrorKind::RedirectLimit, 919 BlossomPhase::Retrieval, 920 false, 921 possible_orphan, 922 attempt, 923 )) 924 } 925 926 #[cfg_attr(coverage_nightly, coverage(off))] 927 async fn hardened_client( 928 config: &BlossomConfig, 929 endpoint: &BlossomEndpoint, 930 cancellation: &BlossomCancellation, 931 phase: BlossomPhase, 932 attempts: u8, 933 possible_orphan: bool, 934 ) -> Result<reqwest::Client, BlossomError> { 935 let addresses = resolve( 936 endpoint, 937 config.connect_timeout(), 938 cancellation, 939 phase, 940 attempts, 941 possible_orphan, 942 ) 943 .await?; 944 hardened_client_with_addresses( 945 config, 946 endpoint, 947 addresses.as_slice(), 948 phase, 949 attempts, 950 possible_orphan, 951 ) 952 } 953 954 fn hardened_client_with_addresses( 955 config: &BlossomConfig, 956 endpoint: &BlossomEndpoint, 957 addresses: &[SocketAddr], 958 phase: BlossomPhase, 959 attempts: u8, 960 possible_orphan: bool, 961 ) -> Result<reqwest::Client, BlossomError> { 962 let mut builder = reqwest::Client::builder() 963 .redirect(reqwest::redirect::Policy::none()) 964 .connect_timeout(config.connect_timeout()) 965 .timeout(config.request_timeout()) 966 .pool_max_idle_per_host(0); 967 if endpoint.host().parse::<std::net::IpAddr>().is_err() { 968 let address = addresses.first().ok_or_else(|| { 969 failure( 970 BlossomErrorKind::ResolutionFailed, 971 phase, 972 true, 973 possible_orphan, 974 attempts, 975 ) 976 })?; 977 builder = builder.resolve(endpoint.host(), *address); 978 } 979 builder.build().map_err(|_| { 980 failure( 981 BlossomErrorKind::Transport, 982 phase, 983 true, 984 possible_orphan, 985 attempts, 986 ) 987 }) 988 } 989 990 #[cfg_attr(coverage_nightly, coverage(off))] 991 async fn resolve( 992 endpoint: &BlossomEndpoint, 993 timeout: Duration, 994 cancellation: &BlossomCancellation, 995 phase: BlossomPhase, 996 attempts: u8, 997 possible_orphan: bool, 998 ) -> Result<Vec<SocketAddr>, BlossomError> { 999 let lookup = tokio::net::lookup_host((endpoint.host(), endpoint.port())); 1000 let mut resolved = tokio::select! { 1001 biased; 1002 _ = cancellation.cancelled() => { 1003 return Err(failure( 1004 BlossomErrorKind::Cancelled, 1005 phase, 1006 true, 1007 possible_orphan, 1008 attempts, 1009 )); 1010 } 1011 result = tokio::time::timeout(timeout, lookup) => result 1012 .map_err(|_| failure( 1013 BlossomErrorKind::Timeout, 1014 phase, 1015 true, 1016 possible_orphan, 1017 attempts, 1018 ))? 1019 .map_err(|_| { 1020 failure( 1021 BlossomErrorKind::ResolutionFailed, 1022 phase, 1023 true, 1024 possible_orphan, 1025 attempts, 1026 ) 1027 })?, 1028 }; 1029 let addresses = resolved 1030 .by_ref() 1031 .take(MAX_RESOLVED_ADDRESSES + 1) 1032 .collect::<Vec<_>>(); 1033 if addresses.is_empty() || addresses.len() > MAX_RESOLVED_ADDRESSES { 1034 return Err(failure( 1035 BlossomErrorKind::ResolutionFailed, 1036 phase, 1037 true, 1038 possible_orphan, 1039 attempts, 1040 )); 1041 } 1042 endpoint 1043 .validate_resolved_addresses(addresses.iter().map(SocketAddr::ip)) 1044 .map_err(|error| with_operation(error, possible_orphan, attempts))?; 1045 Ok(addresses) 1046 } 1047 1048 #[cfg_attr(coverage_nightly, coverage(off))] 1049 async fn read_bounded( 1050 mut response: reqwest::Response, 1051 max_bytes: usize, 1052 cancellation: &BlossomCancellation, 1053 phase: BlossomPhase, 1054 possible_orphan: bool, 1055 attempts: u8, 1056 ) -> Result<Vec<u8>, BlossomError> { 1057 if response 1058 .content_length() 1059 .is_some_and(|length| length > max_bytes as u64) 1060 { 1061 return Err(failure( 1062 BlossomErrorKind::ResponseTooLarge, 1063 phase, 1064 false, 1065 possible_orphan, 1066 attempts, 1067 )); 1068 } 1069 let mut bytes = Vec::with_capacity( 1070 response 1071 .content_length() 1072 .and_then(|length| usize::try_from(length).ok()) 1073 .unwrap_or(0) 1074 .min(max_bytes), 1075 ); 1076 loop { 1077 let pending = response.chunk(); 1078 let chunk = tokio::select! { 1079 biased; 1080 _ = cancellation.cancelled() => { 1081 return Err(failure( 1082 BlossomErrorKind::Cancelled, 1083 phase, 1084 true, 1085 possible_orphan, 1086 attempts, 1087 )); 1088 } 1089 chunk = pending => chunk.map_err(|error| request_error(error, phase, possible_orphan, attempts))?, 1090 }; 1091 let Some(chunk) = chunk else { 1092 break; 1093 }; 1094 if bytes.len().saturating_add(chunk.len()) > max_bytes { 1095 return Err(failure( 1096 BlossomErrorKind::ResponseTooLarge, 1097 phase, 1098 false, 1099 possible_orphan, 1100 attempts, 1101 )); 1102 } 1103 bytes.extend_from_slice(&chunk); 1104 } 1105 Ok(bytes) 1106 } 1107 1108 async fn retry_delay( 1109 config: &BlossomConfig, 1110 attempt: u8, 1111 cancellation: &BlossomCancellation, 1112 phase: BlossomPhase, 1113 possible_orphan: bool, 1114 ) -> Result<(), BlossomError> { 1115 let delay = config.retry_delay(attempt); 1116 tokio::select! { 1117 biased; 1118 _ = cancellation.cancelled() => Err(failure( 1119 BlossomErrorKind::Cancelled, 1120 phase, 1121 true, 1122 possible_orphan, 1123 attempt, 1124 )), 1125 _ = tokio::time::sleep(delay) => Ok(()), 1126 } 1127 } 1128 1129 fn ensure_not_cancelled( 1130 cancellation: &BlossomCancellation, 1131 phase: BlossomPhase, 1132 attempts: u8, 1133 possible_orphan: bool, 1134 ) -> Result<(), BlossomError> { 1135 if cancellation.is_cancelled() { 1136 Err(failure( 1137 BlossomErrorKind::Cancelled, 1138 phase, 1139 true, 1140 possible_orphan, 1141 attempts, 1142 )) 1143 } else { 1144 Ok(()) 1145 } 1146 } 1147 1148 fn request_error( 1149 error: reqwest::Error, 1150 phase: BlossomPhase, 1151 possible_orphan: bool, 1152 attempts: u8, 1153 ) -> BlossomError { 1154 let kind = if error.is_timeout() { 1155 BlossomErrorKind::Timeout 1156 } else { 1157 BlossomErrorKind::Transport 1158 }; 1159 failure(kind, phase, true, possible_orphan, attempts) 1160 } 1161 1162 fn http_status_error( 1163 status: StatusCode, 1164 phase: BlossomPhase, 1165 possible_orphan: bool, 1166 attempts: u8, 1167 ) -> BlossomError { 1168 let retryable = matches!( 1169 status, 1170 StatusCode::REQUEST_TIMEOUT 1171 | StatusCode::TOO_EARLY 1172 | StatusCode::TOO_MANY_REQUESTS 1173 | StatusCode::INTERNAL_SERVER_ERROR 1174 | StatusCode::BAD_GATEWAY 1175 | StatusCode::SERVICE_UNAVAILABLE 1176 | StatusCode::GATEWAY_TIMEOUT 1177 ); 1178 failure( 1179 BlossomErrorKind::HttpStatus, 1180 phase, 1181 retryable, 1182 possible_orphan, 1183 attempts, 1184 ) 1185 .with_http_status(status.as_u16()) 1186 } 1187 1188 async fn http_status_response_error( 1189 response: reqwest::Response, 1190 phase: BlossomPhase, 1191 possible_orphan: bool, 1192 attempts: u8, 1193 cancellation: &BlossomCancellation, 1194 ) -> BlossomError { 1195 let status = response.status(); 1196 let mut error = http_status_error(status, phase, possible_orphan, attempts); 1197 if let Some(code) = read_server_error_code(response, cancellation).await { 1198 error = error.with_server_error_code(code); 1199 } 1200 error 1201 } 1202 1203 async fn read_server_error_code( 1204 mut response: reqwest::Response, 1205 cancellation: &BlossomCancellation, 1206 ) -> Option<String> { 1207 if response 1208 .content_length() 1209 .is_some_and(|length| length > MAX_ERROR_RESPONSE_BYTES as u64) 1210 { 1211 return None; 1212 } 1213 let mut bytes = Vec::with_capacity( 1214 response 1215 .content_length() 1216 .and_then(|length| usize::try_from(length).ok()) 1217 .unwrap_or(0) 1218 .min(MAX_ERROR_RESPONSE_BYTES), 1219 ); 1220 loop { 1221 let chunk = tokio::select! { 1222 biased; 1223 _ = cancellation.cancelled() => return None, 1224 chunk = response.chunk() => chunk.ok()?, 1225 }; 1226 let Some(chunk) = chunk else { 1227 break; 1228 }; 1229 if bytes.len().saturating_add(chunk.len()) > MAX_ERROR_RESPONSE_BYTES { 1230 return None; 1231 } 1232 bytes.extend_from_slice(&chunk); 1233 } 1234 parse_server_error_code(bytes.as_slice()) 1235 } 1236 1237 fn parse_server_error_code(bytes: &[u8]) -> Option<String> { 1238 if bytes.len() > MAX_ERROR_RESPONSE_BYTES { 1239 return None; 1240 } 1241 let value: serde_json::Value = serde_json::from_slice(bytes).ok()?; 1242 let code = value.get("error")?.as_str()?; 1243 if code.is_empty() 1244 || code.len() > MAX_SERVER_ERROR_CODE_BYTES 1245 || !code 1246 .bytes() 1247 .all(|byte| byte.is_ascii_lowercase() || byte.is_ascii_digit() || byte == b'_') 1248 { 1249 return None; 1250 } 1251 Some(code.to_owned()) 1252 } 1253 1254 fn with_operation(error: BlossomError, possible_orphan: bool, attempts: u8) -> BlossomError { 1255 error.with_operation(possible_orphan, attempts) 1256 } 1257 1258 fn failure( 1259 kind: BlossomErrorKind, 1260 phase: BlossomPhase, 1261 retryable: bool, 1262 possible_orphan: bool, 1263 attempts: u8, 1264 ) -> BlossomError { 1265 BlossomError::new(kind, phase, retryable, possible_orphan, attempts) 1266 } 1267 1268 pub(crate) fn verify_image( 1269 bytes: &[u8], 1270 media_type: &MediaType, 1271 expected: BlossomImageDimensions, 1272 ) -> Result<(), BlossomError> { 1273 let dimensions = inspect_image(bytes, media_type)?; 1274 if dimensions != expected { 1275 return Err(failure( 1276 BlossomErrorKind::DimensionMismatch, 1277 BlossomPhase::Verification, 1278 false, 1279 false, 1280 0, 1281 )); 1282 } 1283 Ok(()) 1284 } 1285 1286 fn inspect_image( 1287 bytes: &[u8], 1288 media_type: &MediaType, 1289 ) -> Result<BlossomImageDimensions, BlossomError> { 1290 let detected = detect_image(bytes).ok_or_else(|| { 1291 failure( 1292 BlossomErrorKind::InvalidImageBytes, 1293 BlossomPhase::Verification, 1294 false, 1295 false, 1296 0, 1297 ) 1298 })?; 1299 let declared = match media_type.as_str() { 1300 "image/png" => ImageKind::Png, 1301 "image/jpeg" => ImageKind::Jpeg, 1302 "image/gif" => ImageKind::Gif, 1303 "image/webp" => ImageKind::Webp, 1304 _ => { 1305 return Err(failure( 1306 BlossomErrorKind::UnsupportedMediaType, 1307 BlossomPhase::Verification, 1308 false, 1309 false, 1310 0, 1311 )); 1312 } 1313 }; 1314 if declared != detected.0 { 1315 return Err(failure( 1316 BlossomErrorKind::MediaTypeMismatch, 1317 BlossomPhase::Verification, 1318 false, 1319 false, 1320 0, 1321 )); 1322 } 1323 Ok(detected.1) 1324 } 1325 1326 #[derive(Clone, Copy, Debug, Eq, PartialEq)] 1327 enum ImageKind { 1328 Png, 1329 Jpeg, 1330 Gif, 1331 Webp, 1332 } 1333 1334 fn detect_image(bytes: &[u8]) -> Option<(ImageKind, BlossomImageDimensions)> { 1335 detect_png(bytes) 1336 .map(|dimensions| (ImageKind::Png, dimensions)) 1337 .or_else(|| detect_jpeg(bytes).map(|dimensions| (ImageKind::Jpeg, dimensions))) 1338 .or_else(|| detect_gif(bytes).map(|dimensions| (ImageKind::Gif, dimensions))) 1339 .or_else(|| detect_webp(bytes).map(|dimensions| (ImageKind::Webp, dimensions))) 1340 } 1341 1342 fn detect_png(bytes: &[u8]) -> Option<BlossomImageDimensions> { 1343 if bytes.len() < 24 || &bytes[..8] != b"\x89PNG\r\n\x1a\n" || &bytes[12..16] != b"IHDR" { 1344 return None; 1345 } 1346 BlossomImageDimensions::new( 1347 u32::from_be_bytes(bytes[16..20].try_into().ok()?), 1348 u32::from_be_bytes(bytes[20..24].try_into().ok()?), 1349 ) 1350 .ok() 1351 } 1352 1353 fn detect_gif(bytes: &[u8]) -> Option<BlossomImageDimensions> { 1354 if bytes.len() < 10 || !matches!(&bytes[..6], b"GIF87a" | b"GIF89a") { 1355 return None; 1356 } 1357 BlossomImageDimensions::new( 1358 u32::from(u16::from_le_bytes(bytes[6..8].try_into().ok()?)), 1359 u32::from(u16::from_le_bytes(bytes[8..10].try_into().ok()?)), 1360 ) 1361 .ok() 1362 } 1363 1364 fn detect_jpeg(bytes: &[u8]) -> Option<BlossomImageDimensions> { 1365 if bytes.len() < 4 || bytes[..2] != [0xff, 0xd8] { 1366 return None; 1367 } 1368 let mut offset = 2_usize; 1369 while offset < bytes.len() { 1370 while bytes.get(offset) == Some(&0xff) { 1371 offset += 1; 1372 } 1373 let marker = *bytes.get(offset)?; 1374 offset += 1; 1375 if marker == 0xd9 || marker == 0xda { 1376 return None; 1377 } 1378 if marker == 0x01 || (0xd0..=0xd7).contains(&marker) { 1379 continue; 1380 } 1381 let length = usize::from(u16::from_be_bytes([ 1382 *bytes.get(offset)?, 1383 *bytes.get(offset + 1)?, 1384 ])); 1385 if length < 2 || offset.checked_add(length)? > bytes.len() { 1386 return None; 1387 } 1388 if matches!( 1389 marker, 1390 0xc0 | 0xc1 1391 | 0xc2 1392 | 0xc3 1393 | 0xc5 1394 | 0xc6 1395 | 0xc7 1396 | 0xc9 1397 | 0xca 1398 | 0xcb 1399 | 0xcd 1400 | 0xce 1401 | 0xcf 1402 ) { 1403 if length < 7 { 1404 return None; 1405 } 1406 return BlossomImageDimensions::new( 1407 u32::from(u16::from_be_bytes([ 1408 *bytes.get(offset + 5)?, 1409 *bytes.get(offset + 6)?, 1410 ])), 1411 u32::from(u16::from_be_bytes([ 1412 *bytes.get(offset + 3)?, 1413 *bytes.get(offset + 4)?, 1414 ])), 1415 ) 1416 .ok(); 1417 } 1418 offset += length; 1419 } 1420 None 1421 } 1422 1423 fn detect_webp(bytes: &[u8]) -> Option<BlossomImageDimensions> { 1424 if bytes.len() < 30 || &bytes[..4] != b"RIFF" || &bytes[8..12] != b"WEBP" { 1425 return None; 1426 } 1427 match &bytes[12..16] { 1428 b"VP8X" => BlossomImageDimensions::new( 1429 little_u24(&bytes[24..27])?.checked_add(1)?, 1430 little_u24(&bytes[27..30])?.checked_add(1)?, 1431 ) 1432 .ok(), 1433 b"VP8L" if bytes.get(20) == Some(&0x2f) => { 1434 let packed = u32::from_le_bytes(bytes[21..25].try_into().ok()?); 1435 BlossomImageDimensions::new( 1436 (packed & 0x3fff).checked_add(1)?, 1437 ((packed >> 14) & 0x3fff).checked_add(1)?, 1438 ) 1439 .ok() 1440 } 1441 b"VP8 " if bytes.get(23..26) == Some(&[0x9d, 0x01, 0x2a]) => BlossomImageDimensions::new( 1442 u32::from(u16::from_le_bytes(bytes[26..28].try_into().ok()?) & 0x3fff), 1443 u32::from(u16::from_le_bytes(bytes[28..30].try_into().ok()?) & 0x3fff), 1444 ) 1445 .ok(), 1446 _ => None, 1447 } 1448 } 1449 1450 fn little_u24(bytes: &[u8]) -> Option<u32> { 1451 Some( 1452 u32::from(*bytes.first()?) 1453 | u32::from(*bytes.get(1)?) << 8 1454 | u32::from(*bytes.get(2)?) << 16, 1455 ) 1456 } 1457 1458 #[cfg(test)] 1459 mod native_completion_tests; 1460 1461 #[cfg(test)] 1462 mod tests { 1463 use super::*; 1464 use radroots_blossom::Sha256; 1465 use std::sync::Arc; 1466 use tokio::io::{AsyncReadExt, AsyncWriteExt}; 1467 use tokio::net::TcpListener; 1468 1469 static LOOPBACK_TEST_GUARD: tokio::sync::Mutex<()> = tokio::sync::Mutex::const_new(()); 1470 1471 pub(super) fn png(width: u32, height: u32) -> Vec<u8> { 1472 let mut bytes = b"\x89PNG\r\n\x1a\n\0\0\0\rIHDR".to_vec(); 1473 bytes.extend_from_slice(&width.to_be_bytes()); 1474 bytes.extend_from_slice(&height.to_be_bytes()); 1475 bytes 1476 } 1477 1478 #[test] 1479 fn image_headers_bind_mime_and_dimensions() { 1480 let dimensions = BlossomImageDimensions::new(1200, 900).expect("dimensions"); 1481 let bytes = png(1200, 900); 1482 assert!(verify_image(&bytes, &MediaType::parse("image/png").unwrap(), dimensions).is_ok()); 1483 assert_eq!( 1484 verify_image(&bytes, &MediaType::parse("image/jpeg").unwrap(), dimensions) 1485 .expect_err("wrong MIME") 1486 .kind(), 1487 BlossomErrorKind::MediaTypeMismatch 1488 ); 1489 assert_eq!( 1490 verify_image( 1491 &bytes, 1492 &MediaType::parse("image/png").unwrap(), 1493 BlossomImageDimensions::new(1, 1).unwrap(), 1494 ) 1495 .expect_err("wrong dimensions") 1496 .kind(), 1497 BlossomErrorKind::DimensionMismatch 1498 ); 1499 assert_eq!( 1500 verify_image( 1501 b"not an image", 1502 &MediaType::parse("image/png").unwrap(), 1503 dimensions 1504 ) 1505 .expect_err("invalid image") 1506 .kind(), 1507 BlossomErrorKind::InvalidImageBytes 1508 ); 1509 } 1510 1511 #[test] 1512 fn supported_image_headers_are_bounded_and_nonzero() { 1513 let gif = b"GIF89a\x02\0\x03\0"; 1514 assert_eq!( 1515 detect_gif(gif), 1516 Some(BlossomImageDimensions::new(2, 3).unwrap()) 1517 ); 1518 1519 let jpeg = [ 1520 0xff, 0xd8, 0xff, 0xc0, 0x00, 0x11, 0x08, 0x00, 0x03, 0x00, 0x02, 0x03, 0x01, 0x11, 1521 0x00, 0x02, 0x11, 0x00, 0x03, 0x11, 0x00, 1522 ]; 1523 assert_eq!( 1524 detect_jpeg(&jpeg), 1525 Some(BlossomImageDimensions::new(2, 3).unwrap()) 1526 ); 1527 1528 let mut webp = vec![0_u8; 30]; 1529 webp[..4].copy_from_slice(b"RIFF"); 1530 webp[8..12].copy_from_slice(b"WEBP"); 1531 webp[12..16].copy_from_slice(b"VP8X"); 1532 webp[24..27].copy_from_slice(&[1, 0, 0]); 1533 webp[27..30].copy_from_slice(&[2, 0, 0]); 1534 assert_eq!( 1535 detect_webp(&webp), 1536 Some(BlossomImageDimensions::new(2, 3).unwrap()) 1537 ); 1538 } 1539 1540 #[test] 1541 fn image_header_decoders_reject_every_malformed_boundary() { 1542 let mut invalid_png_signature = vec![0_u8; 24]; 1543 invalid_png_signature[12..16].copy_from_slice(b"IHDR"); 1544 let mut invalid_png_chunk = png(2, 3); 1545 invalid_png_chunk[12..16].copy_from_slice(b"NOPE"); 1546 let zero_png = png(0, 3); 1547 assert_eq!(detect_png(b"short"), None); 1548 assert_eq!(detect_png(&invalid_png_signature), None); 1549 assert_eq!(detect_png(&invalid_png_chunk), None); 1550 assert_eq!(detect_png(&zero_png), None); 1551 1552 assert_eq!(detect_gif(b"short"), None); 1553 assert_eq!(detect_gif(b"GIF00a\x02\0\x03\0"), None); 1554 assert_eq!( 1555 detect_gif(b"GIF87a\x02\0\x03\0"), 1556 Some(BlossomImageDimensions::new(2, 3).unwrap()) 1557 ); 1558 assert_eq!(detect_gif(b"GIF89a\0\0\x03\0"), None); 1559 1560 for malformed in [ 1561 vec![0xff, 0xd8], 1562 vec![0, 0, 0, 0], 1563 vec![0xff, 0xd8, 0xff, 0xd9], 1564 vec![0xff, 0xd8, 0xff, 0xda], 1565 vec![0xff, 0xd8, 0x01], 1566 vec![0xff, 0xd8, 0xd0], 1567 vec![0xff, 0xd8, 0xff, 0xe0, 0, 1], 1568 vec![0xff, 0xd8, 0xff, 0xe0, 0, 100], 1569 vec![0xff, 0xd8, 0xff, 0xc0, 0, 6, 0, 0, 0, 0], 1570 ] { 1571 assert_eq!(detect_jpeg(&malformed), None, "{malformed:?}"); 1572 } 1573 let jpeg_with_prefix = [ 1574 0xff, 0xd8, 0xff, 0xe0, 0, 2, 0xff, 0xc2, 0, 7, 8, 0, 3, 0, 2, 1575 ]; 1576 assert_eq!( 1577 detect_jpeg(&jpeg_with_prefix), 1578 Some(BlossomImageDimensions::new(2, 3).unwrap()) 1579 ); 1580 1581 assert_eq!(detect_webp(b"short"), None); 1582 let mut wrong_riff = vec![0_u8; 30]; 1583 wrong_riff[8..12].copy_from_slice(b"WEBP"); 1584 assert_eq!(detect_webp(&wrong_riff), None); 1585 let mut wrong_webp = vec![0_u8; 30]; 1586 wrong_webp[..4].copy_from_slice(b"RIFF"); 1587 assert_eq!(detect_webp(&wrong_webp), None); 1588 1589 let mut lossless = vec![0_u8; 30]; 1590 lossless[..4].copy_from_slice(b"RIFF"); 1591 lossless[8..12].copy_from_slice(b"WEBP"); 1592 lossless[12..16].copy_from_slice(b"VP8L"); 1593 lossless[20] = 0x2f; 1594 let packed = 1_u32 | (2_u32 << 14); 1595 lossless[21..25].copy_from_slice(&packed.to_le_bytes()); 1596 assert_eq!( 1597 detect_webp(&lossless), 1598 Some(BlossomImageDimensions::new(2, 3).unwrap()) 1599 ); 1600 lossless[20] = 0; 1601 assert_eq!(detect_webp(&lossless), None); 1602 1603 let mut lossy = vec![0_u8; 30]; 1604 lossy[..4].copy_from_slice(b"RIFF"); 1605 lossy[8..12].copy_from_slice(b"WEBP"); 1606 lossy[12..16].copy_from_slice(b"VP8 "); 1607 lossy[23..26].copy_from_slice(&[0x9d, 0x01, 0x2a]); 1608 lossy[26..28].copy_from_slice(&2_u16.to_le_bytes()); 1609 lossy[28..30].copy_from_slice(&3_u16.to_le_bytes()); 1610 assert_eq!( 1611 detect_webp(&lossy), 1612 Some(BlossomImageDimensions::new(2, 3).unwrap()) 1613 ); 1614 lossy[23] = 0; 1615 assert_eq!(detect_webp(&lossy), None); 1616 lossy[12..16].copy_from_slice(b"NOPE"); 1617 assert_eq!(detect_webp(&lossy), None); 1618 1619 assert_eq!(little_u24(&[]), None); 1620 assert_eq!(little_u24(&[1]), None); 1621 assert_eq!(little_u24(&[1, 2]), None); 1622 assert_eq!(little_u24(&[1, 2, 3]), Some(0x03_02_01)); 1623 } 1624 1625 #[test] 1626 fn image_verification_supports_every_declared_mime() { 1627 let gif = b"GIF89a\x02\0\x03\0"; 1628 assert!( 1629 verify_image( 1630 gif, 1631 &MediaType::parse("image/gif").unwrap(), 1632 BlossomImageDimensions::new(2, 3).unwrap(), 1633 ) 1634 .is_ok() 1635 ); 1636 1637 let jpeg = [0xff, 0xd8, 0xff, 0xc0, 0, 7, 8, 0, 3, 0, 2]; 1638 assert!( 1639 verify_image( 1640 &jpeg, 1641 &MediaType::parse("image/jpeg").unwrap(), 1642 BlossomImageDimensions::new(2, 3).unwrap(), 1643 ) 1644 .is_ok() 1645 ); 1646 1647 let mut webp = vec![0_u8; 30]; 1648 webp[..4].copy_from_slice(b"RIFF"); 1649 webp[8..12].copy_from_slice(b"WEBP"); 1650 webp[12..16].copy_from_slice(b"VP8X"); 1651 webp[24..27].copy_from_slice(&[1, 0, 0]); 1652 webp[27..30].copy_from_slice(&[2, 0, 0]); 1653 assert!( 1654 verify_image( 1655 &webp, 1656 &MediaType::parse("image/webp").unwrap(), 1657 BlossomImageDimensions::new(2, 3).unwrap(), 1658 ) 1659 .is_ok() 1660 ); 1661 assert_eq!( 1662 verify_image( 1663 &webp, 1664 &MediaType::parse("application/octet-stream").unwrap(), 1665 BlossomImageDimensions::new(2, 3).unwrap(), 1666 ) 1667 .expect_err("unsupported media type") 1668 .kind(), 1669 BlossomErrorKind::UnsupportedMediaType 1670 ); 1671 } 1672 1673 #[tokio::test] 1674 async fn retry_and_error_classification_are_bounded_and_redacted() { 1675 for status in [ 1676 StatusCode::REQUEST_TIMEOUT, 1677 StatusCode::TOO_EARLY, 1678 StatusCode::TOO_MANY_REQUESTS, 1679 StatusCode::INTERNAL_SERVER_ERROR, 1680 StatusCode::BAD_GATEWAY, 1681 StatusCode::SERVICE_UNAVAILABLE, 1682 StatusCode::GATEWAY_TIMEOUT, 1683 ] { 1684 assert!(http_status_error(status, BlossomPhase::Upload, true, 1).retryable()); 1685 } 1686 assert!( 1687 !http_status_error(StatusCode::BAD_REQUEST, BlossomPhase::Upload, false, 1).retryable() 1688 ); 1689 1690 let cancellation = BlossomCancellation::default(); 1691 assert!(ensure_not_cancelled(&cancellation, BlossomPhase::Upload, 0, false).is_ok()); 1692 cancellation.cancel(); 1693 assert_eq!( 1694 ensure_not_cancelled(&cancellation, BlossomPhase::Retrieval, 2, true) 1695 .expect_err("cancelled") 1696 .kind(), 1697 BlossomErrorKind::Cancelled 1698 ); 1699 1700 let profile = simulator_profile("http://127.0.0.1:9"); 1701 let config = BlossomConfig::from_profile(profile) 1702 .with_network_policy( 1703 Duration::from_millis(10), 1704 Duration::from_millis(10), 1705 1, 1706 Duration::from_millis(1), 1707 ) 1708 .unwrap(); 1709 let delay = BlossomCancellation::default(); 1710 assert!( 1711 retry_delay(&config, 1, &delay, BlossomPhase::Upload, false) 1712 .await 1713 .is_ok() 1714 ); 1715 delay.cancel(); 1716 assert_eq!( 1717 retry_delay(&config, 20, &delay, BlossomPhase::Retrieval, true,) 1718 .await 1719 .expect_err("cancelled delay") 1720 .kind(), 1721 BlossomErrorKind::Cancelled 1722 ); 1723 1724 let transport_error = reqwest::Client::new() 1725 .get("http://127.0.0.1:9") 1726 .send() 1727 .await 1728 .expect_err("closed port"); 1729 assert_eq!( 1730 request_error(transport_error, BlossomPhase::Upload, false, 1).kind(), 1731 BlossomErrorKind::Transport 1732 ); 1733 1734 let bytes = png(2, 3); 1735 let request = upload_request("http://127.0.0.1:9", bytes.clone()); 1736 let too_small = BlossomConfig::from_profile(simulator_profile("http://127.0.0.1:9")) 1737 .with_limits(1, 100, 0) 1738 .unwrap(); 1739 assert_eq!( 1740 upload_with_authorization( 1741 too_small, 1742 request, 1743 "Nostr redacted", 1744 BlossomCancellation::default(), 1745 ) 1746 .await 1747 .expect_err("oversize request") 1748 .kind(), 1749 BlossomErrorKind::ResponseTooLarge 1750 ); 1751 } 1752 1753 #[test] 1754 fn server_error_codes_are_bounded_validated_and_detail_free() { 1755 assert_eq!( 1756 parse_server_error_code( 1757 br#"{"error":"entitlement_missing","detail":"must never be retained"}"# 1758 ) 1759 .as_deref(), 1760 Some("entitlement_missing") 1761 ); 1762 assert_eq!( 1763 parse_server_error_code(br#"{"error":"unsafe-value"}"#), 1764 None 1765 ); 1766 assert_eq!(parse_server_error_code(br#"{"error":"UPPERCASE"}"#), None); 1767 assert_eq!( 1768 parse_server_error_code(format!(r#"{{"error":"{}"}}"#, "a".repeat(65)).as_bytes()), 1769 None 1770 ); 1771 assert_eq!( 1772 parse_server_error_code(vec![b' '; MAX_ERROR_RESPONSE_BYTES + 1].as_slice()), 1773 None 1774 ); 1775 assert_eq!(parse_server_error_code(b"not-json"), None); 1776 } 1777 1778 #[tokio::test] 1779 async fn http_failure_retains_only_the_public_server_error_code() { 1780 let _loopback_guard = LOOPBACK_TEST_GUARD.lock().await; 1781 let listener = TcpListener::bind("127.0.0.1:0").await.expect("bind"); 1782 let address = listener.local_addr().expect("address"); 1783 let server = tokio::spawn(async move { 1784 let (mut stream, _) = listener.accept().await.expect("accept"); 1785 let request = read_request(&mut stream).await; 1786 assert!(String::from_utf8_lossy(&request).starts_with("PUT /upload HTTP/1.1")); 1787 let body = 1788 br#"{"error":"entitlement_missing","detail":"sensitive operational context"}"#; 1789 let response = format!( 1790 "HTTP/1.1 403 Forbidden\r\nContent-Type: application/json\r\nContent-Length: {}\r\nConnection: close\r\n\r\n", 1791 body.len() 1792 ); 1793 stream.write_all(response.as_bytes()).await.expect("head"); 1794 stream.write_all(body).await.expect("body"); 1795 stream.shutdown().await.expect("close"); 1796 }); 1797 let origin = format!("http://{address}"); 1798 let config = config(origin.as_str()); 1799 let endpoint = config.profile().primary().clone(); 1800 let error = upload_once( 1801 &config, 1802 &endpoint, 1803 &upload_request(origin.as_str(), png(2, 3)), 1804 "Nostr redacted", 1805 &BlossomCancellation::default(), 1806 1, 1807 ) 1808 .await 1809 .expect_err("forbidden upload"); 1810 assert_eq!(error.http_status(), Some(403)); 1811 assert_eq!(error.server_error_code(), Some("entitlement_missing")); 1812 assert!(!format!("{error:?}").contains("sensitive operational context")); 1813 assert!(!format!("{error:?}").contains("Nostr redacted")); 1814 server.await.expect("server"); 1815 } 1816 1817 #[test] 1818 fn descriptor_verification_checks_each_identity_field() { 1819 let bytes = png(2, 3); 1820 let request = upload_request("http://127.0.0.1:3000", bytes); 1821 let expected_url = 1822 BlobUrl::parse(format!("http://127.0.0.1:3000/{}.png", request.sha256()).as_str()) 1823 .unwrap(); 1824 let descriptor = |url: BlobUrl, size: u64, media_type: &str| { 1825 BlobDescriptor::new( 1826 url, 1827 request.sha256(), 1828 size, 1829 MediaType::parse(media_type).unwrap(), 1830 1, 1831 ) 1832 .unwrap() 1833 }; 1834 assert!( 1835 verify_descriptor( 1836 &request, 1837 &expected_url, 1838 descriptor( 1839 BlobUrl::parse( 1840 format!("http://localhost:3000/{}.png", request.sha256()).as_str() 1841 ) 1842 .unwrap(), 1843 request.byte_size(), 1844 "image/png", 1845 ), 1846 1, 1847 ) 1848 .is_err() 1849 ); 1850 assert!( 1851 verify_descriptor( 1852 &request, 1853 &expected_url, 1854 descriptor(expected_url.clone(), request.byte_size() + 1, "image/png"), 1855 1, 1856 ) 1857 .is_err() 1858 ); 1859 assert!( 1860 verify_descriptor( 1861 &request, 1862 &expected_url, 1863 descriptor(expected_url.clone(), request.byte_size(), "image/jpeg",), 1864 1, 1865 ) 1866 .is_err() 1867 ); 1868 } 1869 1870 #[derive(Clone, Copy)] 1871 enum RetrievalResponse { 1872 Exact, 1873 Altered, 1874 RedirectExternal, 1875 WrongMime, 1876 Oversize, 1877 Stall, 1878 } 1879 1880 pub(super) async fn read_request(stream: &mut tokio::net::TcpStream) -> Vec<u8> { 1881 let mut request = Vec::new(); 1882 let header_end = loop { 1883 let mut chunk = [0_u8; 1024]; 1884 let read = stream.read(&mut chunk).await.expect("request read"); 1885 assert_ne!(read, 0, "request ended before headers"); 1886 request.extend_from_slice(&chunk[..read]); 1887 if let Some(end) = request.windows(4).position(|value| value == b"\r\n\r\n") { 1888 break end + 4; 1889 } 1890 assert!(request.len() < 64 * 1024, "request headers are bounded"); 1891 }; 1892 let headers = String::from_utf8_lossy(&request[..header_end]); 1893 let content_length = headers 1894 .lines() 1895 .find_map(|line| { 1896 line.to_ascii_lowercase() 1897 .strip_prefix("content-length: ") 1898 .and_then(|value| value.trim().parse::<usize>().ok()) 1899 }) 1900 .unwrap_or(0); 1901 while request.len() - header_end < content_length { 1902 let mut chunk = [0_u8; 1024]; 1903 let read = stream.read(&mut chunk).await.expect("body read"); 1904 assert_ne!(read, 0, "request ended before body"); 1905 request.extend_from_slice(&chunk[..read]); 1906 } 1907 request 1908 } 1909 1910 async fn spawn_server( 1911 bytes: Vec<u8>, 1912 retrieval: RetrievalResponse, 1913 bad_descriptor: bool, 1914 ) -> (String, tokio::task::JoinHandle<Vec<u8>>) { 1915 let listener = TcpListener::bind("127.0.0.1:0").await.expect("bind"); 1916 let address = listener.local_addr().expect("address"); 1917 let origin = format!("http://{address}"); 1918 let expected_hash = Sha256::digest(bytes.as_slice()); 1919 let descriptor_hash = if bad_descriptor { 1920 Sha256::digest(b"different") 1921 } else { 1922 expected_hash 1923 }; 1924 let descriptor_url = format!("{origin}/{descriptor_hash}.png"); 1925 let descriptor = BlobDescriptor::new( 1926 BlobUrl::parse(descriptor_url.as_str()).expect("url"), 1927 descriptor_hash, 1928 bytes.len() as u64, 1929 MediaType::parse("image/png").expect("media type"), 1930 1_900_000_000, 1931 ) 1932 .expect("descriptor"); 1933 let descriptor_json = serde_json::to_vec(&descriptor).expect("descriptor json"); 1934 let task = tokio::spawn(async move { 1935 let (mut upload, _) = listener.accept().await.expect("upload accept"); 1936 let upload_request = read_request(&mut upload).await; 1937 let response = format!( 1938 "HTTP/1.1 200 OK\r\nContent-Type: application/json\r\nContent-Length: {}\r\nConnection: close\r\n\r\n", 1939 descriptor_json.len() 1940 ); 1941 upload 1942 .write_all(response.as_bytes()) 1943 .await 1944 .expect("upload head"); 1945 upload 1946 .write_all(&descriptor_json) 1947 .await 1948 .expect("upload body"); 1949 upload.shutdown().await.expect("upload close"); 1950 1951 if bad_descriptor { 1952 return upload_request; 1953 } 1954 let (mut retrieval_stream, _) = listener.accept().await.expect("retrieval accept"); 1955 let _ = read_request(&mut retrieval_stream).await; 1956 match retrieval { 1957 RetrievalResponse::Exact | RetrievalResponse::Altered => { 1958 let body = if matches!(retrieval, RetrievalResponse::Altered) { 1959 let mut altered = bytes.clone(); 1960 altered[0] ^= 1; 1961 altered 1962 } else { 1963 bytes 1964 }; 1965 let response = format!( 1966 "HTTP/1.1 200 OK\r\nContent-Type: image/png\r\nContent-Length: {}\r\nConnection: close\r\n\r\n", 1967 body.len() 1968 ); 1969 retrieval_stream 1970 .write_all(response.as_bytes()) 1971 .await 1972 .expect("retrieval head"); 1973 retrieval_stream 1974 .write_all(&body) 1975 .await 1976 .expect("retrieval body"); 1977 } 1978 RetrievalResponse::RedirectExternal => { 1979 let location = format!("https://example.com/{expected_hash}.png"); 1980 let response = format!( 1981 "HTTP/1.1 302 Found\r\nLocation: {location}\r\nContent-Length: 0\r\nConnection: close\r\n\r\n" 1982 ); 1983 retrieval_stream 1984 .write_all(response.as_bytes()) 1985 .await 1986 .expect("redirect"); 1987 } 1988 RetrievalResponse::WrongMime => { 1989 let response = format!( 1990 "HTTP/1.1 200 OK\r\nContent-Type: image/jpeg\r\nContent-Length: {}\r\nConnection: close\r\n\r\n", 1991 bytes.len() 1992 ); 1993 retrieval_stream 1994 .write_all(response.as_bytes()) 1995 .await 1996 .expect("wrong MIME"); 1997 retrieval_stream 1998 .write_all(&bytes) 1999 .await 2000 .expect("retrieval body"); 2001 } 2002 RetrievalResponse::Oversize => { 2003 retrieval_stream 2004 .write_all(b"HTTP/1.1 200 OK\r\nContent-Type: image/png\r\nContent-Length: 999999\r\nConnection: close\r\n\r\n") 2005 .await 2006 .expect("oversize"); 2007 } 2008 RetrievalResponse::Stall => { 2009 tokio::time::sleep(Duration::from_secs(3)).await; 2010 } 2011 } 2012 retrieval_stream.shutdown().await.expect("retrieval close"); 2013 upload_request 2014 }); 2015 (origin, task) 2016 } 2017 2018 async fn spawn_retry_server(bytes: Vec<u8>) -> (String, tokio::task::JoinHandle<()>) { 2019 let listener = TcpListener::bind("127.0.0.1:0").await.expect("bind"); 2020 let address = listener.local_addr().expect("address"); 2021 let origin = format!("http://{address}"); 2022 let hash = Sha256::digest(bytes.as_slice()); 2023 let descriptor = BlobDescriptor::new( 2024 BlobUrl::parse(format!("{origin}/{hash}.png").as_str()).expect("url"), 2025 hash, 2026 bytes.len() as u64, 2027 MediaType::parse("image/png").expect("media type"), 2028 1_900_000_000, 2029 ) 2030 .expect("descriptor"); 2031 let descriptor_json = serde_json::to_vec(&descriptor).expect("descriptor JSON"); 2032 let task = tokio::spawn(async move { 2033 for step in 0..4 { 2034 let (mut stream, _) = listener.accept().await.expect("accept"); 2035 let _ = read_request(&mut stream).await; 2036 match step { 2037 0 => { 2038 stream 2039 .write_all(b"HTTP/1.1 503 Service Unavailable\r\nContent-Length: 0\r\nConnection: close\r\n\r\n") 2040 .await 2041 .expect("upload retry"); 2042 } 2043 1 => { 2044 let head = format!( 2045 "HTTP/1.1 200 OK\r\nContent-Type: application/json\r\nContent-Length: {}\r\nConnection: close\r\n\r\n", 2046 descriptor_json.len() 2047 ); 2048 stream 2049 .write_all(head.as_bytes()) 2050 .await 2051 .expect("descriptor head"); 2052 stream 2053 .write_all(&descriptor_json) 2054 .await 2055 .expect("descriptor body"); 2056 } 2057 2 => { 2058 stream 2059 .write_all(b"HTTP/1.1 503 Service Unavailable\r\nContent-Length: 0\r\nConnection: close\r\n\r\n") 2060 .await 2061 .expect("retrieval retry"); 2062 } 2063 3 => { 2064 let head = format!( 2065 "HTTP/1.1 200 OK\r\nContent-Type: image/png\r\nContent-Length: {}\r\nConnection: close\r\n\r\n", 2066 bytes.len() 2067 ); 2068 stream 2069 .write_all(head.as_bytes()) 2070 .await 2071 .expect("retrieval head"); 2072 stream.write_all(&bytes).await.expect("retrieval body"); 2073 } 2074 _ => unreachable!(), 2075 } 2076 stream.shutdown().await.expect("close"); 2077 } 2078 }); 2079 (origin, task) 2080 } 2081 2082 pub(super) fn upload_request(_origin: &str, bytes: Vec<u8>) -> BlossomUploadRequest { 2083 BlossomUploadRequest::new( 2084 Arc::from(bytes), 2085 MediaType::parse("image/png").expect("media type"), 2086 BlossomImageDimensions::new(2, 3).expect("dimensions"), 2087 1_900_000_000_000, 2088 ) 2089 .expect("upload request") 2090 } 2091 2092 pub(super) fn config(origin: &str) -> BlossomConfig { 2093 BlossomConfig::from_profile(simulator_profile(origin)) 2094 .with_network_policy( 2095 Duration::from_secs(2), 2096 Duration::from_secs(2), 2097 1, 2098 Duration::from_millis(1), 2099 ) 2100 .expect("network policy") 2101 } 2102 2103 fn simulator_profile(origin: &str) -> crate::transport::BlossomProfile { 2104 crate::transport::BlossomProfile::new( 2105 crate::transport::BlossomHostKind::Simulator, 2106 crate::transport::BlossomEndpointAuthority::LoopbackDevelopment, 2107 origin, 2108 std::iter::empty::<&str>(), 2109 ) 2110 .expect("profile") 2111 } 2112 2113 #[tokio::test] 2114 async fn native_upload_completion_verifies_descriptor_and_retrieved_bytes() { 2115 let bytes = png(2, 3); 2116 let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); 2117 let origin = format!("http://{}", listener.local_addr().unwrap()); 2118 let request = upload_request(origin.as_str(), bytes.clone()); 2119 let expected_url = 2120 BlobUrl::parse(format!("{origin}/{}.png", request.sha256()).as_str()).unwrap(); 2121 let descriptor = BlobDescriptor::new( 2122 expected_url, 2123 request.sha256(), 2124 request.byte_size(), 2125 MediaType::parse("image/png").unwrap(), 2126 1_900_000_000, 2127 ) 2128 .unwrap(); 2129 let descriptor_body = serde_json::to_vec(&descriptor).unwrap(); 2130 let server = tokio::spawn(async move { 2131 let (mut stream, _) = listener.accept().await.unwrap(); 2132 let request = read_request(&mut stream).await; 2133 assert!(String::from_utf8_lossy(&request).starts_with("GET /")); 2134 let head = format!( 2135 "HTTP/1.1 200 OK\r\nContent-Type: image/png\r\nContent-Length: {}\r\nConnection: close\r\n\r\n", 2136 bytes.len() 2137 ); 2138 stream.write_all(head.as_bytes()).await.unwrap(); 2139 stream.write_all(&bytes).await.unwrap(); 2140 stream.shutdown().await.unwrap(); 2141 }); 2142 let slot = crate::transport::BlossomSlot::new(); 2143 slot.configure(config(origin.as_str())).unwrap(); 2144 let transaction = slot 2145 .prepare_upload(upload_request(origin.as_str(), png(2, 3))) 2146 .unwrap(); 2147 let receipt = slot 2148 .complete_native_upload( 2149 transaction, 2150 200, 2151 Some("application/json"), 2152 None, 2153 descriptor_body.as_slice(), 2154 BlossomCancellation::default(), 2155 ) 2156 .await 2157 .unwrap(); 2158 assert_eq!(receipt.descriptor().sha256(), descriptor.sha256()); 2159 assert_eq!(receipt.attempts(), 2); 2160 server.await.unwrap(); 2161 } 2162 2163 async fn spawn_inbound_server( 2164 bytes: Vec<u8>, 2165 content_type: &'static str, 2166 content_encoding: Option<&'static str>, 2167 declared_length: Option<usize>, 2168 ) -> (String, tokio::task::JoinHandle<Vec<u8>>) { 2169 let listener = TcpListener::bind("127.0.0.1:0").await.expect("bind"); 2170 let address = listener.local_addr().expect("address"); 2171 let origin = format!("http://{address}"); 2172 let task = tokio::spawn(async move { 2173 let (mut stream, _) = listener.accept().await.expect("accept"); 2174 let request = read_request(&mut stream).await; 2175 let encoding = content_encoding 2176 .map(|value| format!("Content-Encoding: {value}\r\n")) 2177 .unwrap_or_default(); 2178 let response = format!( 2179 "HTTP/1.1 200 OK\r\nContent-Type: {content_type}\r\n{encoding}Content-Length: {}\r\nConnection: close\r\n\r\n", 2180 declared_length.unwrap_or(bytes.len()) 2181 ); 2182 stream.write_all(response.as_bytes()).await.expect("head"); 2183 stream.write_all(&bytes).await.expect("body"); 2184 stream.shutdown().await.expect("close"); 2185 request 2186 }); 2187 (origin, task) 2188 } 2189 2190 fn inbound_request( 2191 origin: &str, 2192 bytes: &[u8], 2193 dimensions: BlossomImageDimensions, 2194 ) -> BlossomInboundRequest { 2195 let hash = Sha256::digest(bytes); 2196 BlossomInboundRequest::new( 2197 BlobUrl::parse(format!("{origin}/{hash}.png").as_str()).expect("URL"), 2198 Some(MediaType::parse("image/png").expect("media type")), 2199 Some(bytes.len() as u64), 2200 Some(dimensions), 2201 ) 2202 .expect("inbound request") 2203 } 2204 2205 #[tokio::test] 2206 async fn inbound_retrieval_returns_only_exact_verified_image_bytes() { 2207 let bytes = png(2, 3); 2208 let (origin, server) = spawn_inbound_server(bytes.clone(), "image/png", None, None).await; 2209 let slot = crate::transport::BlossomSlot::new(); 2210 slot.configure(config(origin.as_str())).expect("configure"); 2211 let receipt = slot 2212 .retrieve( 2213 inbound_request( 2214 origin.as_str(), 2215 bytes.as_slice(), 2216 BlossomImageDimensions::new(2, 3).unwrap(), 2217 ), 2218 BlossomCancellation::default(), 2219 ) 2220 .await 2221 .expect("inbound receipt"); 2222 assert_eq!(receipt.bytes(), bytes.as_slice()); 2223 assert_eq!(receipt.commitment().sha256(), Sha256::digest(&bytes)); 2224 assert_eq!(receipt.commitment().size(), bytes.len() as u64); 2225 assert_eq!(receipt.commitment().media_type().as_str(), "image/png"); 2226 assert_eq!( 2227 receipt.dimensions(), 2228 BlossomImageDimensions::new(2, 3).unwrap() 2229 ); 2230 assert_eq!(receipt.attempts(), 1); 2231 assert_eq!( 2232 receipt.config_fingerprint(), 2233 slot.config_fingerprint().expect("fingerprint") 2234 ); 2235 let request = server.await.expect("server"); 2236 let headers = String::from_utf8_lossy(&request); 2237 assert!(headers.starts_with("GET /")); 2238 assert!( 2239 headers 2240 .to_ascii_lowercase() 2241 .contains("accept-encoding: identity") 2242 ); 2243 assert!(!headers.to_ascii_lowercase().contains("authorization:")); 2244 } 2245 2246 #[tokio::test] 2247 async fn inbound_retrieval_rejects_encoding_truncation_hash_and_dimensions() { 2248 let exact = png(2, 3); 2249 let cases = [ 2250 ( 2251 exact.clone(), 2252 Some("gzip"), 2253 None, 2254 BlossomImageDimensions::new(2, 3).unwrap(), 2255 BlossomErrorKind::ContentEncodingDenied, 2256 ), 2257 ( 2258 exact[..exact.len() - 1].to_vec(), 2259 None, 2260 Some(exact.len()), 2261 BlossomImageDimensions::new(2, 3).unwrap(), 2262 BlossomErrorKind::Transport, 2263 ), 2264 ( 2265 { 2266 let mut altered = exact.clone(); 2267 let last = altered.len() - 1; 2268 altered[last] ^= 1; 2269 altered 2270 }, 2271 None, 2272 None, 2273 BlossomImageDimensions::new(2, 3).unwrap(), 2274 BlossomErrorKind::ResponseHashMismatch, 2275 ), 2276 ( 2277 exact.clone(), 2278 None, 2279 None, 2280 BlossomImageDimensions::new(3, 2).unwrap(), 2281 BlossomErrorKind::DimensionMismatch, 2282 ), 2283 ]; 2284 for (served, encoding, length, dimensions, expected) in cases { 2285 let (origin, server) = 2286 spawn_inbound_server(served, "image/png", encoding, length).await; 2287 let slot = crate::transport::BlossomSlot::new(); 2288 slot.configure(config(origin.as_str())).expect("configure"); 2289 let error = slot 2290 .retrieve( 2291 inbound_request(origin.as_str(), exact.as_slice(), dimensions), 2292 BlossomCancellation::default(), 2293 ) 2294 .await 2295 .expect_err("hostile response"); 2296 assert_eq!(error.kind(), expected); 2297 assert!(!error.possible_orphan()); 2298 server.await.expect("server"); 2299 } 2300 } 2301 2302 #[tokio::test] 2303 async fn inbound_retry_and_cancellation_are_bounded_and_operation_scoped() { 2304 let bytes = png(2, 3); 2305 let listener = TcpListener::bind("127.0.0.1:0").await.expect("bind"); 2306 let origin = format!("http://{}", listener.local_addr().expect("address")); 2307 let server_bytes = bytes.clone(); 2308 let server = tokio::spawn(async move { 2309 for attempt in 0..2 { 2310 let (mut stream, _) = listener.accept().await.expect("accept"); 2311 let _ = read_request(&mut stream).await; 2312 if attempt == 0 { 2313 stream 2314 .write_all(b"HTTP/1.1 503 Service Unavailable\r\nContent-Length: 0\r\nConnection: close\r\n\r\n") 2315 .await 2316 .expect("retry response"); 2317 } else { 2318 let response = format!( 2319 "HTTP/1.1 200 OK\r\nContent-Type: image/png\r\nContent-Length: {}\r\nConnection: close\r\n\r\n", 2320 server_bytes.len() 2321 ); 2322 stream.write_all(response.as_bytes()).await.expect("head"); 2323 stream.write_all(&server_bytes).await.expect("body"); 2324 } 2325 stream.shutdown().await.expect("close"); 2326 } 2327 }); 2328 let slot = crate::transport::BlossomSlot::new(); 2329 slot.configure( 2330 BlossomConfig::from_profile(simulator_profile(origin.as_str())) 2331 .with_network_policy( 2332 Duration::from_millis(100), 2333 Duration::from_millis(100), 2334 2, 2335 Duration::from_millis(1), 2336 ) 2337 .unwrap(), 2338 ) 2339 .unwrap(); 2340 let receipt = slot 2341 .retrieve( 2342 inbound_request( 2343 origin.as_str(), 2344 bytes.as_slice(), 2345 BlossomImageDimensions::new(2, 3).unwrap(), 2346 ), 2347 BlossomCancellation::default(), 2348 ) 2349 .await 2350 .expect("retried retrieval"); 2351 assert_eq!(receipt.attempts(), 2); 2352 server.await.expect("server"); 2353 2354 let cancellation = BlossomCancellation::default(); 2355 cancellation.cancel(); 2356 let error = slot 2357 .retrieve( 2358 inbound_request( 2359 origin.as_str(), 2360 bytes.as_slice(), 2361 BlossomImageDimensions::new(2, 3).unwrap(), 2362 ), 2363 cancellation, 2364 ) 2365 .await 2366 .expect_err("cancelled retrieval"); 2367 assert_eq!(error.kind(), BlossomErrorKind::Cancelled); 2368 assert_eq!(error.phase(), BlossomPhase::Retrieval); 2369 assert!(!error.possible_orphan()); 2370 } 2371 2372 #[tokio::test] 2373 async fn non_mutating_probe_records_only_dns_transport_and_http_evidence() { 2374 let _loopback_guard = LOOPBACK_TEST_GUARD.lock().await; 2375 let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); 2376 let address = listener.local_addr().unwrap(); 2377 let server = tokio::spawn(async move { 2378 let (mut stream, _) = listener.accept().await.unwrap(); 2379 let mut request = vec![0_u8; 2_048]; 2380 let count = stream.read(&mut request).await.unwrap(); 2381 let request = String::from_utf8_lossy(&request[..count]); 2382 assert!(request.starts_with(&format!("GET /{BLOSSOM_PROBE_HASH} HTTP/1.1"))); 2383 assert!(!request.to_ascii_lowercase().contains("authorization:")); 2384 stream 2385 .write_all(b"HTTP/1.1 404 Not Found\r\nContent-Length: 0\r\n\r\n") 2386 .await 2387 .unwrap(); 2388 }); 2389 let slot = crate::transport::BlossomSlot::new(); 2390 slot.configure(config(format!("http://{address}").as_str())) 2391 .unwrap(); 2392 let initial = slot.evidence().unwrap(); 2393 assert_eq!( 2394 initial.state(), 2395 crate::transport::BlossomEvidenceState::ConfiguredUnobserved 2396 ); 2397 assert!(initial.observed_at_unix_ms().is_none()); 2398 2399 let observed = slot.probe(BlossomCancellation::default()).await.unwrap(); 2400 assert_eq!( 2401 observed.state(), 2402 crate::transport::BlossomEvidenceState::TlsHttpObserved 2403 ); 2404 assert_eq!(observed.http_status(), Some(404)); 2405 assert!(observed.error_code().is_none()); 2406 assert!(observed.observed_at_unix_ms().is_some()); 2407 server.await.unwrap(); 2408 } 2409 2410 #[tokio::test] 2411 async fn cancelled_probe_is_retryable_redacted_and_never_claims_dns() { 2412 let slot = crate::transport::BlossomSlot::new(); 2413 slot.configure(config("http://127.0.0.1:9")).unwrap(); 2414 let cancellation = BlossomCancellation::default(); 2415 cancellation.cancel(); 2416 let error = slot.probe(cancellation).await.unwrap_err(); 2417 assert_eq!(error.kind(), BlossomErrorKind::Cancelled); 2418 let evidence = slot.evidence().unwrap(); 2419 assert_eq!( 2420 evidence.state(), 2421 crate::transport::BlossomEvidenceState::RetryableFailure 2422 ); 2423 assert_eq!( 2424 evidence.last_successful_state(), 2425 crate::transport::BlossomEvidenceState::ConfiguredUnobserved 2426 ); 2427 assert_eq!(evidence.error_code(), Some("blossom_cancelled")); 2428 assert_eq!(evidence.error_phase(), Some(BlossomPhase::Probe)); 2429 assert!(!evidence.possible_orphan()); 2430 } 2431 2432 #[tokio::test] 2433 async fn reconfiguration_during_probe_cannot_promote_stale_evidence() { 2434 let _loopback_guard = LOOPBACK_TEST_GUARD.lock().await; 2435 let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); 2436 let address = listener.local_addr().unwrap(); 2437 let (accepted_tx, accepted_rx) = tokio::sync::oneshot::channel(); 2438 let (release_tx, release_rx) = tokio::sync::oneshot::channel(); 2439 let server = tokio::spawn(async move { 2440 let (mut stream, _) = listener.accept().await.unwrap(); 2441 let mut request = vec![0_u8; 2_048]; 2442 let _ = stream.read(&mut request).await.unwrap(); 2443 accepted_tx.send(()).unwrap(); 2444 release_rx.await.unwrap(); 2445 stream 2446 .write_all(b"HTTP/1.1 404 Not Found\r\nContent-Length: 0\r\n\r\n") 2447 .await 2448 .unwrap(); 2449 }); 2450 let slot = crate::transport::BlossomSlot::new(); 2451 slot.configure(config(format!("http://{address}").as_str())) 2452 .unwrap(); 2453 let probing = { 2454 let slot = slot.clone(); 2455 tokio::spawn(async move { slot.probe(BlossomCancellation::default()).await }) 2456 }; 2457 accepted_rx.await.unwrap(); 2458 slot.configure(config("http://127.0.0.1:9")).unwrap(); 2459 release_tx.send(()).unwrap(); 2460 let error = probing.await.unwrap().unwrap_err(); 2461 assert_eq!(error.kind(), BlossomErrorKind::ConfigurationChanged); 2462 assert_eq!(error.phase(), BlossomPhase::Probe); 2463 assert!(!error.possible_orphan()); 2464 let evidence = slot.evidence().unwrap(); 2465 assert_eq!(evidence.origin(), "http://127.0.0.1:9"); 2466 assert_eq!( 2467 evidence.state(), 2468 crate::transport::BlossomEvidenceState::ConfiguredUnobserved 2469 ); 2470 server.await.unwrap(); 2471 } 2472 2473 #[tokio::test] 2474 async fn loopback_upload_preserves_exact_bytes_and_verifies_retrieval() { 2475 let _loopback_guard = LOOPBACK_TEST_GUARD.lock().await; 2476 let bytes = png(2, 3); 2477 let (origin, server) = spawn_server(bytes.clone(), RetrievalResponse::Exact, false).await; 2478 let receipt = upload_with_authorization( 2479 config(origin.as_str()), 2480 upload_request(origin.as_str(), bytes.clone()), 2481 "Nostr secret-token-value", 2482 BlossomCancellation::default(), 2483 ) 2484 .await 2485 .expect("verified upload"); 2486 assert_eq!(receipt.descriptor().size(), bytes.len() as u64); 2487 assert_eq!( 2488 receipt.dimensions(), 2489 BlossomImageDimensions::new(2, 3).unwrap() 2490 ); 2491 let upload_wire = server.await.expect("server"); 2492 let header_end = upload_wire 2493 .windows(4) 2494 .position(|value| value == b"\r\n\r\n") 2495 .unwrap() 2496 + 4; 2497 assert_eq!(&upload_wire[header_end..], bytes.as_slice()); 2498 let headers = String::from_utf8_lossy(&upload_wire[..header_end]); 2499 assert!(headers.contains("authorization: Nostr secret-token-value")); 2500 assert!(!format!("{receipt:?}").contains("secret-token-value")); 2501 } 2502 2503 #[tokio::test] 2504 async fn retryable_upload_and_retrieval_failures_recover_with_bounded_attempts() { 2505 let _loopback_guard = LOOPBACK_TEST_GUARD.lock().await; 2506 let bytes = png(2, 3); 2507 let (origin, server) = spawn_retry_server(bytes.clone()).await; 2508 let config = BlossomConfig::from_profile(simulator_profile(origin.as_str())) 2509 .with_network_policy( 2510 Duration::from_secs(2), 2511 Duration::from_secs(2), 2512 2, 2513 Duration::from_millis(1), 2514 ) 2515 .expect("network policy"); 2516 let receipt = upload_with_authorization( 2517 config, 2518 upload_request(origin.as_str(), bytes), 2519 "Nostr redacted", 2520 BlossomCancellation::default(), 2521 ) 2522 .await 2523 .expect("retry recovery"); 2524 assert_eq!(receipt.attempts(), 4); 2525 server.await.expect("server"); 2526 } 2527 2528 #[tokio::test] 2529 async fn descriptor_redirect_body_mime_timeout_and_cancellation_fail_closed() { 2530 let _loopback_guard = LOOPBACK_TEST_GUARD.lock().await; 2531 let cases = [ 2532 ( 2533 RetrievalResponse::Exact, 2534 true, 2535 BlossomErrorKind::DescriptorMismatch, 2536 ), 2537 ( 2538 RetrievalResponse::RedirectExternal, 2539 false, 2540 BlossomErrorKind::UnsafeRedirect, 2541 ), 2542 ( 2543 RetrievalResponse::Altered, 2544 false, 2545 BlossomErrorKind::ResponseHashMismatch, 2546 ), 2547 ( 2548 RetrievalResponse::WrongMime, 2549 false, 2550 BlossomErrorKind::MediaTypeMismatch, 2551 ), 2552 ( 2553 RetrievalResponse::Oversize, 2554 false, 2555 BlossomErrorKind::ResponseSizeMismatch, 2556 ), 2557 (RetrievalResponse::Stall, false, BlossomErrorKind::Timeout), 2558 ]; 2559 for (response, bad_descriptor, expected) in cases { 2560 let bytes = png(2, 3); 2561 let (origin, server) = spawn_server(bytes.clone(), response, bad_descriptor).await; 2562 let error = upload_with_authorization( 2563 config(origin.as_str()), 2564 upload_request(origin.as_str(), bytes), 2565 "Nostr redacted", 2566 BlossomCancellation::default(), 2567 ) 2568 .await 2569 .expect_err("must fail closed"); 2570 assert_eq!(error.kind(), expected); 2571 assert!(error.possible_orphan()); 2572 assert!(!format!("{error:?}").contains("Nostr redacted")); 2573 server.await.expect("server"); 2574 } 2575 2576 let cancellation = BlossomCancellation::default(); 2577 cancellation.cancel(); 2578 let bytes = png(2, 3); 2579 let request = BlossomUploadRequest::new( 2580 Arc::from(bytes), 2581 MediaType::parse("image/png").unwrap(), 2582 BlossomImageDimensions::new(2, 3).unwrap(), 2583 1, 2584 ) 2585 .unwrap(); 2586 let error = upload_with_authorization( 2587 config("http://127.0.0.1:9"), 2588 request, 2589 "Nostr redacted", 2590 cancellation, 2591 ) 2592 .await 2593 .expect_err("cancelled"); 2594 assert_eq!(error.kind(), BlossomErrorKind::Cancelled); 2595 assert!(!error.possible_orphan()); 2596 } 2597 }