lib

Core libraries for Radroots
git clone https://radroots.dev/git/lib.git
Log | Files | Refs | README

commit 462c662b83eb8d2ace6a3aae496aeba32dbc31d3
parent 20ec51df367ca50658f94740d06b27a22ad8ae69
Author: triesap <tyson@radroots.org>
Date:   Fri,  7 Aug 2026 19:11:00 +0000

feat(sdk): add verified Blossom transport

- enforce bounded BUD-11 authorization and exact-byte BUD-02 uploads
- validate BUD-01 retrievals, redirects, descriptors, media, and network policy
- persist verified media, retry evidence, cancellation, and possible orphans
- prove simulator I/O, redaction, package boundaries, and release coverage

Diffstat:
MCargo.lock | 29+++++++++++++++++++++++++++++
Mcrates/mobile_core/src/runtime/builder.rs | 35+++++++++++++++++++++++++++++++++++
Mcrates/mobile_core/src/runtime/mod.rs | 13+++++++++++++
Mcrates/mobile_core/src/runtime/product_surface.rs | 5+++--
Mcrates/mobile_core/src/runtime/product_surface/outbox.rs | 429++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-----------
Mcrates/mobile_core/src/runtime/sdk.rs | 61+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcrates/sdk/Cargo.toml | 22++++++++++++++++++++--
Acrates/sdk/src/adapters/blossom.rs | 1611+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcrates/sdk/src/adapters/mod.rs | 2++
Mcrates/sdk/src/client.rs | 60++++++++++++++++++++++++++++++++++++++++++++++++++----------
Mcrates/sdk/src/error.rs | 4++--
Mcrates/sdk/src/lib.rs | 1+
Mcrates/sdk/src/transport.rs | 1505++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-
Mcrates/sdk/tests/package_boundary.rs | 3+++
Mcrates/sdk/tests/public_api.rs | 15++++++++++++++-
15 files changed, 3717 insertions(+), 78 deletions(-)

diff --git a/Cargo.lock b/Cargo.lock @@ -4427,6 +4427,7 @@ dependencies = [ "base64", "bytes", "futures-core", + "futures-util", "http", "http-body", "http-body-util", @@ -4446,12 +4447,14 @@ dependencies = [ "sync_wrapper", "tokio", "tokio-rustls", + "tokio-util", "tower", "tower-http", "tower-service", "url", "wasm-bindgen", "wasm-bindgen-futures", + "wasm-streams", "web-sys", "webpki-roots 1.0.6", ] @@ -5479,6 +5482,19 @@ dependencies = [ ] [[package]] +name = "tokio-util" +version = "0.7.19" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "494815d09bf52b5548659851081238f0ca39ff638363907596da739561c62c52" +dependencies = [ + "bytes", + "futures-core", + "futures-sink", + "pin-project-lite", + "tokio", +] + +[[package]] name = "toml" version = "0.5.11" source = "registry+https://github.com/rust-lang/crates.io-index" @@ -6312,6 +6328,19 @@ dependencies = [ ] [[package]] +name = "wasm-streams" +version = "0.4.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "15053d8d85c7eccdbefef60f06769760a563c7f0a9d6902a13d35c7800b0ad65" +dependencies = [ + "futures-util", + "js-sys", + "wasm-bindgen", + "wasm-bindgen-futures", + "web-sys", +] + +[[package]] name = "wasmparser" version = "0.244.0" source = "registry+https://github.com/rust-lang/crates.io-index" diff --git a/crates/mobile_core/src/runtime/builder.rs b/crates/mobile_core/src/runtime/builder.rs @@ -8,6 +8,8 @@ pub struct RuntimeBuilder { signer: Option<std::sync::Arc<dyn radroots_signing::Signer>>, #[cfg(feature = "mobile-social")] relay_profile: radroots_sdk::transport::RelayProfile, + #[cfg(feature = "mobile-social")] + blossom_config: Option<radroots_sdk::transport::BlossomConfig>, } impl RuntimeBuilder { @@ -20,6 +22,8 @@ impl RuntimeBuilder { #[cfg(feature = "mobile-social")] relay_profile: radroots_sdk::transport::RelayProfile::public(Vec::<String>::new()) .expect("bundled public relay profile is valid"), + #[cfg(feature = "mobile-social")] + blossom_config: None, } } @@ -40,6 +44,17 @@ impl RuntimeBuilder { self } + /// Installs one validated inert Blossom environment profile. + #[cfg(feature = "mobile-social")] + #[must_use] + pub fn blossom_config( + mut self, + blossom_config: radroots_sdk::transport::BlossomConfig, + ) -> Self { + self.blossom_config = Some(blossom_config); + self + } + /// Opens the exact authenticated user's durable SQLite store. pub async fn build(self) -> Result<RadrootsRuntime, RadrootsAppError> { if self.store.protected_data() == ProtectedDataAvailability::Unavailable { @@ -57,6 +72,8 @@ impl RuntimeBuilder { self.signer, #[cfg(feature = "mobile-social")] Some(self.relay_profile), + #[cfg(feature = "mobile-social")] + self.blossom_config, ) } } @@ -146,6 +163,24 @@ mod tests { .expect("configured profile"), report ); + assert!( + runtime + .configure_simulator_blossom(vec!["http://127.0.0.1:3000".to_owned()]) + .is_ok() + ); + assert_eq!( + runtime.sdk_blossom_profile().expect("Blossom profile"), + Some("simulator_local".to_owned()) + ); + assert!( + runtime + .configure_public_blossom(vec!["http://127.0.0.1:3000".to_owned()]) + .is_err() + ); + assert_eq!( + runtime.sdk_blossom_profile().expect("unchanged profile"), + Some("simulator_local".to_owned()) + ); runtime.shutdown().await.expect("shutdown"); } diff --git a/crates/mobile_core/src/runtime/mod.rs b/crates/mobile_core/src/runtime/mod.rs @@ -37,6 +37,9 @@ impl RadrootsRuntime { #[cfg(feature = "mobile-social")] relay_profile: Option< radroots_sdk::transport::RelayProfile, >, + #[cfg(feature = "mobile-social")] blossom_config: Option< + radroots_sdk::transport::BlossomConfig, + >, ) -> Result<Self, RadrootsAppError> { #[cfg(feature = "mobile-social")] let builder = { @@ -48,6 +51,14 @@ impl RadrootsRuntime { } let builder = builder .nostr(nostr_slot) + .blossom({ + let slot = radroots_sdk::transport::BlossomSlot::new(); + if let Some(config) = blossom_config { + slot.configure(config) + .map_err(|error| RadrootsAppError::runtime(error.code().to_owned()))?; + } + slot + }) .host_sync(radroots_sdk::sync::HostPolicy::standard()); match signer { Some(signer) => builder.signing(radroots_sdk::signing::Provider::host(signer)), @@ -74,6 +85,8 @@ impl RadrootsRuntime { None, #[cfg(feature = "mobile-social")] None, + #[cfg(feature = "mobile-social")] + None, ) } diff --git a/crates/mobile_core/src/runtime/product_surface.rs b/crates/mobile_core/src/runtime/product_surface.rs @@ -34,8 +34,9 @@ pub use model::{ }; #[cfg(feature = "mobile-social")] pub use outbox::{ - Phase1CancellationPolicy, Phase1DraftError, Phase1DraftStatus, Phase1MediaPrerequisite, - Phase1MediaStage, Phase1OutboxState, Phase1QueuePolicy, Phase1RelaySatisfaction, + Phase1CancellationPolicy, Phase1DraftError, Phase1DraftStatus, Phase1MediaOrphanRecord, + Phase1MediaPrerequisite, Phase1MediaStage, Phase1OutboxState, Phase1QueuePolicy, + Phase1RelaySatisfaction, }; pub use projection::{ProductEventClassification, ProductEventExclusion, classify_admitted_event}; pub use ranking::{RankError, TODAY_RANK_SCHEMA_VERSION, TimeRelevance, TodayRank, TodayRankInput}; diff --git a/crates/mobile_core/src/runtime/product_surface/outbox.rs b/crates/mobile_core/src/runtime/product_surface/outbox.rs @@ -1,6 +1,8 @@ use std::collections::BTreeSet; -use radroots_blossom::{BlobUrl, MediaType, authorization::AuthoredUploadClaim}; +use radroots_blossom::{ + BlobUrl, ByteVerifiedDescriptor, MediaType, authorization::AuthoredUploadClaim, +}; use radroots_event::contract::AuthorRole; use radroots_event_codec::authoring::PlanWireV1; use radroots_identity::PublicKey; @@ -64,37 +66,50 @@ pub struct Phase1MediaPrerequisite { byte_size: u64, stage: Phase1MediaStage, failure_code: Option<String>, + upload_attempts: u8, + verified_at_unix_ms: Option<u64>, + orphan: Option<Phase1MediaOrphanRecord>, +} + +/// Durable, secret-safe evidence that a remote blob may be unreferenced. +#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)] +#[serde(deny_unknown_fields)] +pub struct Phase1MediaOrphanRecord { + reason_code: String, + recorded_at_unix_ms: u64, +} + +impl Phase1MediaOrphanRecord { + pub fn reason_code(&self) -> &str { + self.reason_code.as_str() + } + + pub const fn recorded_at_unix_ms(&self) -> u64 { + self.recorded_at_unix_ms + } } impl Phase1MediaPrerequisite { pub fn new( local_reference: impl Into<String>, - url: impl Into<String>, - sha256: impl Into<String>, - media_type: impl Into<String>, - byte_size: u64, - stage: Phase1MediaStage, + descriptor: &ByteVerifiedDescriptor, ) -> Result<Self, Phase1DraftError> { let value = Self { local_reference: local_reference.into(), - url: url.into(), - sha256: sha256.into(), - media_type: media_type.into(), - byte_size, - stage, + url: descriptor.url().as_blob_url().as_str().to_owned(), + sha256: descriptor.sha256().to_hex(), + media_type: descriptor.media_type().as_str().to_owned(), + byte_size: descriptor.size(), + stage: Phase1MediaStage::Pending, failure_code: None, + upload_attempts: 0, + verified_at_unix_ms: None, + orphan: None, }; value.validate()?; Ok(value) } - pub fn with_failure_code(mut self, code: impl Into<String>) -> Result<Self, Phase1DraftError> { - self.stage = Phase1MediaStage::Failed; - self.failure_code = Some(code.into()); - self.validate()?; - Ok(self) - } - fn validate(&self) -> Result<(), Phase1DraftError> { let blob = BlobUrl::parse(self.url.as_str()).map_err(|_| Phase1DraftError::InvalidMedia)?; let hash = blob.hash_path().hash().to_string(); @@ -111,7 +126,36 @@ impl Phase1MediaPrerequisite { || code != code.trim() || code.chars().any(char::is_control) }) - || (self.stage == Phase1MediaStage::Failed) != self.failure_code.is_some() + || self.orphan.as_ref().is_some_and(|record| { + record.reason_code.is_empty() + || record.reason_code.len() > DRAFT_FAILURE_CODE_MAX_BYTES + || record.reason_code != record.reason_code.trim() + || record.reason_code.chars().any(char::is_control) + || record.recorded_at_unix_ms == 0 + }) + || match self.stage { + Phase1MediaStage::Pending | Phase1MediaStage::Preparing => { + self.failure_code.is_some() + || self.upload_attempts != 0 + || self.verified_at_unix_ms.is_some() + || self.orphan.is_some() + } + Phase1MediaStage::Uploading => { + self.failure_code.is_some() + || self.verified_at_unix_ms.is_some() + || self.orphan.is_some() + } + Phase1MediaStage::Verified => { + self.failure_code.is_some() + || self.upload_attempts == 0 + || self.verified_at_unix_ms.is_none() + || self.orphan.is_some() + } + Phase1MediaStage::Failed => { + self.failure_code.is_none() || self.verified_at_unix_ms.is_some() + } + Phase1MediaStage::Orphaned => self.failure_code.is_some() || self.orphan.is_none(), + } { return Err(Phase1DraftError::InvalidMedia); } @@ -124,6 +168,34 @@ impl Phase1MediaPrerequisite { pub const fn stage(&self) -> Phase1MediaStage { self.stage } + + pub const fn upload_attempts(&self) -> u8 { + self.upload_attempts + } + + pub const fn verified_at_unix_ms(&self) -> Option<u64> { + self.verified_at_unix_ms + } + + pub const fn orphan(&self) -> Option<&Phase1MediaOrphanRecord> { + self.orphan.as_ref() + } + + fn matches_receipt(&self, receipt: &radroots_sdk::transport::BlossomUploadReceipt) -> bool { + let descriptor = receipt.descriptor(); + self.url == descriptor.url().as_blob_url().as_str() + && self.sha256 == descriptor.sha256().to_hex() + && self.media_type == descriptor.media_type().as_str() + && self.byte_size == descriptor.size() + && receipt.attempts() > 0 + && receipt.verified_at_unix_ms() > 0 + } + + fn is_remote_verified(&self) -> bool { + self.stage == Phase1MediaStage::Verified + && self.upload_attempts > 0 + && self.verified_at_unix_ms.is_some() + } } /// Closed delivery-satisfaction profiles available to Phase 1 Add. @@ -505,11 +577,19 @@ impl RadrootsRuntime { .iter_mut() .find(|media| media.url == url) .ok_or(Phase1DraftError::InvalidMedia)?; - if !valid_media_transition(media.stage, stage) { + if !valid_media_transition(media.stage, stage) + || matches!( + stage, + Phase1MediaStage::Verified | Phase1MediaStage::Orphaned + ) + { return Err(Phase1DraftError::InvalidMedia); } media.stage = stage; media.failure_code = failure_code; + if stage != Phase1MediaStage::Failed { + media.failure_code = None; + } media.validate()?; let next_stage = match stage { Phase1MediaStage::Pending | Phase1MediaStage::Preparing => { @@ -530,6 +610,134 @@ impl RadrootsRuntime { self.draft_status_from(receipt.draft().clone()).await } + /// Records the only proof that can advance media to remote-byte-verified. + pub async fn phase1_complete_draft_media( + &self, + draft_id: [u8; 16], + expected_revision: u64, + url: &str, + receipt: radroots_sdk::transport::BlossomUploadReceipt, + updated_at_unix_ms: u64, + ) -> Result<Phase1DraftStatus, Phase1DraftError> { + let draft_id = + AuthoredDraftId::new(draft_id).map_err(|_| Phase1DraftError::InvalidDraft)?; + let expected = AuthoredDraftRevision::new(expected_revision) + .map_err(|_| Phase1DraftError::RevisionConflict)?; + let storage = self.storage()?; + let head = storage + .authored_draft_head(draft_id) + .await + .map_err(|_| Phase1DraftError::Storage)? + .ok_or(Phase1DraftError::NotFound)?; + if head.revision() != expected + || head.stage().is_terminal() + || matches!( + head.stage(), + AuthoredDraftStage::ReadyToSign | AuthoredDraftStage::Queued + ) + { + return Err(Phase1DraftError::RevisionConflict); + } + let mut payload = Phase1DraftPayload::decode(&head)?; + let media = payload + .media + .iter_mut() + .find(|media| media.url == url) + .ok_or(Phase1DraftError::InvalidMedia)?; + if media.stage != Phase1MediaStage::Uploading || !media.matches_receipt(&receipt) { + return Err(Phase1DraftError::InvalidMedia); + } + media.stage = Phase1MediaStage::Verified; + media.failure_code = None; + media.upload_attempts = receipt.attempts(); + media.verified_at_unix_ms = Some(receipt.verified_at_unix_ms()); + media.orphan = None; + media.validate()?; + let next = head + .successor( + payload.encode()?, + AuthoredDraftStage::MediaUploading, + None, + updated_at_unix_ms, + ) + .map_err(|_| Phase1DraftError::RevisionConflict)?; + let receipt = storage + .append_authored_draft(next, Some(expected)) + .await + .map_err(map_draft_storage_error)?; + self.draft_status_from(receipt.draft().clone()).await + } + + /// Persists a redacted recoverable Blossom failure and possible-orphan evidence. + pub async fn phase1_fail_draft_media( + &self, + draft_id: [u8; 16], + expected_revision: u64, + url: &str, + error: &radroots_sdk::transport::BlossomError, + updated_at_unix_ms: u64, + ) -> Result<Phase1DraftStatus, Phase1DraftError> { + if updated_at_unix_ms == 0 { + return Err(Phase1DraftError::InvalidMedia); + } + let draft_id = + AuthoredDraftId::new(draft_id).map_err(|_| Phase1DraftError::InvalidDraft)?; + let expected = AuthoredDraftRevision::new(expected_revision) + .map_err(|_| Phase1DraftError::RevisionConflict)?; + let storage = self.storage()?; + let head = storage + .authored_draft_head(draft_id) + .await + .map_err(|_| Phase1DraftError::Storage)? + .ok_or(Phase1DraftError::NotFound)?; + if head.revision() != expected + || head.stage().is_terminal() + || matches!( + head.stage(), + AuthoredDraftStage::ReadyToSign | AuthoredDraftStage::Queued + ) + { + return Err(Phase1DraftError::RevisionConflict); + } + let mut payload = Phase1DraftPayload::decode(&head)?; + let media = payload + .media + .iter_mut() + .find(|media| media.url == url) + .ok_or(Phase1DraftError::InvalidMedia)?; + if !matches!( + media.stage, + Phase1MediaStage::Pending + | Phase1MediaStage::Preparing + | Phase1MediaStage::Uploading + | Phase1MediaStage::Failed + ) { + return Err(Phase1DraftError::InvalidMedia); + } + media.stage = Phase1MediaStage::Failed; + media.failure_code = Some(error.code().to_owned()); + media.upload_attempts = error.attempts(); + media.verified_at_unix_ms = None; + media.orphan = error.possible_orphan().then(|| Phase1MediaOrphanRecord { + reason_code: error.code().to_owned(), + recorded_at_unix_ms: updated_at_unix_ms, + }); + media.validate()?; + let next = head + .successor( + payload.encode()?, + AuthoredDraftStage::MediaUploading, + None, + updated_at_unix_ms, + ) + .map_err(|_| Phase1DraftError::RevisionConflict)?; + let receipt = storage + .append_authored_draft(next, Some(expected)) + .await + .map_err(map_draft_storage_error)?; + self.draft_status_from(receipt.draft().clone()).await + } + /// Freezes queue intent and atomically prepares the canonical outbox before /// returning `queued`. Network connectivity is neither read nor required. pub async fn phase1_queue_draft( @@ -573,7 +781,7 @@ impl RadrootsRuntime { if payload .media .iter() - .any(|media| media.stage != Phase1MediaStage::Verified) + .any(|media| !media.is_remote_verified()) { return Err(Phase1DraftError::MediaNotReady); } @@ -700,6 +908,115 @@ impl RadrootsRuntime { .map_err(|_| Phase1DraftError::Operation) } + /// Runs the complete durable BUD-11/BUD-02/BUD-01 media transaction. + /// + /// Final image bytes are bound before authorization. The draft becomes + /// verified only after the upload descriptor and a full retrieval agree. + #[allow(clippy::too_many_arguments)] + pub async fn phase1_upload_draft_media( + &self, + draft_id: [u8; 16], + expected_revision: u64, + request: radroots_sdk::transport::BlossomUploadRequest, + authorization_content: radroots_blossom::authorization::AuthorizationContent, + authorization_created_at_unix_s: u64, + authorization_lifetime_seconds: u64, + operation_id: [u8; 16], + artifact_id: [u8; 16], + signing_deadline_unix_ms: u64, + signing_cancellation: Phase1CancellationPolicy, + transfer_cancellation: radroots_sdk::transport::BlossomCancellation, + updated_at_unix_ms: u64, + ) -> Result<Phase1DraftStatus, Phase1DraftError> { + let url = request.expected_url().as_str().to_owned(); + let uploading = self + .phase1_update_draft_media( + draft_id, + expected_revision, + url.as_str(), + Phase1MediaStage::Uploading, + None, + updated_at_unix_ms, + ) + .await?; + let revision = uploading.draft().revision().get(); + let blossom = self + .client + .blossom() + .map_err(|_| Phase1DraftError::OperationUnavailable)? + .ok_or(Phase1DraftError::OperationUnavailable)?; + let claim = match blossom.authored_upload_claim( + &request, + authorization_content, + authorization_created_at_unix_s, + authorization_lifetime_seconds, + ) { + Ok(claim) => claim, + Err(_) => { + self.phase1_update_draft_media( + draft_id, + revision, + url.as_str(), + Phase1MediaStage::Failed, + Some("blossom_authorization_failed".to_owned()), + updated_at_unix_ms, + ) + .await?; + return Err(Phase1DraftError::Operation); + } + }; + let authorization = match self + .phase1_authorize_blossom_upload( + operation_id, + artifact_id, + claim, + signing_deadline_unix_ms, + signing_cancellation, + ) + .await + { + Ok(authorization) => authorization, + Err(error) => { + self.phase1_update_draft_media( + draft_id, + revision, + url.as_str(), + Phase1MediaStage::Failed, + Some("blossom_authorization_failed".to_owned()), + updated_at_unix_ms, + ) + .await?; + return Err(error); + } + }; + match blossom + .upload(request, authorization, transfer_cancellation) + .await + { + Ok(receipt) => { + self.phase1_complete_draft_media( + draft_id, + revision, + url.as_str(), + receipt, + updated_at_unix_ms, + ) + .await + } + Err(error) => { + self.phase1_fail_draft_media( + draft_id, + revision, + url.as_str(), + &error, + updated_at_unix_ms, + ) + .await?; + Err(Phase1DraftError::Operation) + } + } + } + /// Returns durable draft state composed with canonical authored-operation state. pub async fn phase1_draft_status( &self, @@ -765,10 +1082,10 @@ impl RadrootsRuntime { .await .map_err(|_| Phase1DraftError::Operation)?; if status.artifact().signed().is_none() { - mark_possible_orphans(&mut payload.media); + mark_possible_orphans(&mut payload.media, cancelled_at_unix_ms); } } else { - mark_possible_orphans(&mut payload.media); + mark_possible_orphans(&mut payload.media, cancelled_at_unix_ms); } let next = head .successor( @@ -990,10 +1307,15 @@ fn draft_stage_for_media(media: &[Phase1MediaPrerequisite]) -> AuthoredDraftStag } } -fn mark_possible_orphans(media: &mut [Phase1MediaPrerequisite]) { +fn mark_possible_orphans(media: &mut [Phase1MediaPrerequisite], recorded_at_unix_ms: u64) { for media in media { - if media.stage == Phase1MediaStage::Verified { + if media.stage == Phase1MediaStage::Verified || media.orphan.is_some() { media.stage = Phase1MediaStage::Orphaned; + media.failure_code = None; + media.orphan = Some(Phase1MediaOrphanRecord { + reason_code: "draft_cancelled_after_upload".to_owned(), + recorded_at_unix_ms, + }); } } } @@ -1180,6 +1502,7 @@ mod tests { Some(PublicKey::from_hex(AUTHOR).unwrap()), None, None, + None, ) .unwrap() } @@ -1194,6 +1517,7 @@ mod tests { Some(PublicKey::from_hex(AUTHOR).unwrap()), Some(std::sync::Arc::new(signer)), None, + None, ) .unwrap() } @@ -1437,11 +1761,13 @@ mod tests { } #[tokio::test] - async fn media_phase_revisions_gate_queue_and_record_possible_orphans() { + async fn media_phase_revisions_gate_queue_and_reject_forged_verification() { let runtime = runtime(); let id = [9; 16]; let (command, mut media) = photo_command(); media.stage = Phase1MediaStage::Pending; + media.upload_attempts = 0; + media.verified_at_unix_ms = None; media.validate().unwrap(); let saved = runtime .phase1_save_draft(id, command, 1_784_347_200, vec![media], None, 30) @@ -1477,41 +1803,27 @@ mod tests { ) .await .unwrap(); - let verified = runtime - .phase1_update_draft_media( - id, - 3, - uploading.media()[0].url(), - Phase1MediaStage::Verified, - None, - 33, - ) - .await - .unwrap(); - assert_eq!(verified.media()[0].stage(), Phase1MediaStage::Verified); assert_eq!( runtime .phase1_update_draft_media( id, - 4, - verified.media()[0].url(), - Phase1MediaStage::Uploading, + 3, + uploading.media()[0].url(), + Phase1MediaStage::Verified, None, - 34, + 33, ) .await .unwrap_err(), Phase1DraftError::InvalidMedia ); - let queued = runtime - .phase1_queue_draft(id, 4, policy(), 34) - .await - .unwrap(); - let cancelled = runtime - .phase1_cancel_draft(id, queued.draft().revision().get(), 35) - .await - .unwrap(); - assert_eq!(cancelled.media()[0].stage(), Phase1MediaStage::Orphaned); + assert_eq!( + runtime + .phase1_queue_draft(id, 3, policy(), 34) + .await + .unwrap_err(), + Phase1DraftError::MediaNotReady + ); assert_eq!(runtime.phase1_draft_heads(10).await.unwrap().len(), 1); } @@ -1554,6 +1866,8 @@ mod tests { .unwrap() .verify_bytes(bytes, &media_type) .unwrap(); + let mut prerequisite = + Phase1MediaPrerequisite::new("protected://draft/photo-1", &descriptor).unwrap(); let image = AuthoredPostImage::new( AuthoredImage::try_from(descriptor).unwrap(), PostImageDimensions::new(1200, 900).unwrap(), @@ -1563,15 +1877,10 @@ mod tests { let command = Phase1AddCommand::CreatePhotoUpdate( CreatePhotoUpdate::new(format!("Harvest photo {url}"), vec![image]).unwrap(), ); - let prerequisite = Phase1MediaPrerequisite::new( - "protected://draft/photo-1", - url, - hash.to_hex(), - "image/webp", - bytes.len() as u64, - Phase1MediaStage::Verified, - ) - .unwrap(); + prerequisite.stage = Phase1MediaStage::Verified; + prerequisite.upload_attempts = 2; + prerequisite.verified_at_unix_ms = Some(1_784_347_100_000); + prerequisite.validate().unwrap(); (command, prerequisite) } diff --git a/crates/mobile_core/src/runtime/sdk.rs b/crates/mobile_core/src/runtime/sdk.rs @@ -122,6 +122,67 @@ impl RadrootsRuntime { .map_err(RadrootsAppError::from_sdk) } + /// Installs explicit public TLS Blossom origins without probing them. + #[cfg(feature = "mobile-social")] + pub fn configure_public_blossom(&self, origins: Vec<String>) -> Result<(), RadrootsAppError> { + self.configure_blossom_profile( + radroots_sdk::transport::BlossomProfile::public(origins) + .map_err(|error| RadrootsAppError::runtime(error.code().to_owned()))?, + ) + } + + /// Installs exact-loopback simulator Blossom origins without probing them. + #[cfg(feature = "mobile-social")] + pub fn configure_simulator_blossom( + &self, + origins: Vec<String>, + ) -> Result<(), RadrootsAppError> { + self.configure_blossom_profile( + radroots_sdk::transport::BlossomProfile::simulator(origins) + .map_err(|error| RadrootsAppError::runtime(error.code().to_owned()))?, + ) + } + + /// Installs explicit physical-device TLS Blossom origins without probing them. + #[cfg(feature = "mobile-social")] + pub fn configure_device_blossom(&self, origins: Vec<String>) -> Result<(), RadrootsAppError> { + self.configure_blossom_profile( + radroots_sdk::transport::BlossomProfile::device(origins) + .map_err(|error| RadrootsAppError::runtime(error.code().to_owned()))?, + ) + } + + #[cfg(feature = "mobile-social")] + fn configure_blossom_profile( + &self, + profile: radroots_sdk::transport::BlossomProfile, + ) -> Result<(), RadrootsAppError> { + self.client + .configure_blossom(radroots_sdk::transport::BlossomConfig::from_profile( + profile, + )) + .map_err(RadrootsAppError::from_sdk) + } + + /// Returns the inert Blossom environment profile, when configured. + #[cfg(feature = "mobile-social")] + pub fn sdk_blossom_profile(&self) -> Result<Option<String>, RadrootsAppError> { + let profile = self + .client + .blossom() + .map_err(RadrootsAppError::from_sdk)? + .and_then(radroots_sdk::transport::BlossomSlot::profile_kind); + Ok(profile.map(|profile| { + match profile { + radroots_sdk::transport::BlossomProfileKind::Public => "public", + radroots_sdk::transport::BlossomProfileKind::Simulator => "simulator_local", + radroots_sdk::transport::BlossomProfileKind::Device => "device_development", + _ => "unknown", + } + .to_owned() + })) + } + /// Returns passive relay evidence without DNS, socket, or probe work. #[cfg(feature = "mobile-social")] pub fn sdk_relay_status(&self) -> Result<Option<SdkRelayStatusReportRecord>, RadrootsAppError> { diff --git a/crates/sdk/Cargo.toml b/crates/sdk/Cargo.toml @@ -47,7 +47,14 @@ memory = ["radroots_storage/memory"] sqlite = ["dep:radroots_storage_sqlite"] sync = ["dep:radroots_sync", "dep:uuid"] nostr = ["sync", "blossom", "dep:radroots_transport_nostr"] -blossom = ["dep:radroots_nostr", "radroots_nostr/blossom"] +blossom = [ + "dep:radroots_blossom", + "dep:radroots_nostr", + "dep:reqwest", + "dep:serde_json", + "dep:tokio", + "radroots_nostr/blossom", +] nip46 = ["nostr", "dep:radroots_nostr_connect"] local-signing = [ "dep:radroots_nostr", @@ -80,6 +87,10 @@ radroots_trade = { workspace = true, default-features = false, features = [ radroots_transport = { workspace = true, default-features = false } radroots_geonames = { workspace = true, optional = true, default-features = false } +radroots_blossom = { workspace = true, optional = true, default-features = false, features = [ + "serde", + "std", +] } radroots_nostr = { workspace = true, optional = true, default-features = false } radroots_nostr_connect = { workspace = true, optional = true, default-features = false } radroots_secrets = { workspace = true, optional = true, default-features = false } @@ -90,17 +101,24 @@ radroots_transport_nostr = { workspace = true, optional = true, default-features reqwest = { workspace = true, optional = true, default-features = false, features = [ "json", "rustls-tls", + "stream", ] } serde = { workspace = true, optional = true, features = ["derive"] } serde_json = { workspace = true, optional = true, features = ["std"] } uuid = { workspace = true, optional = true, features = ["v4"] } +tokio = { workspace = true, optional = true, features = [ + "macros", + "net", + "sync", + "time", +] } [dev-dependencies] nostr = { workspace = true, features = ["std"] } radroots_blossom = { workspace = true } serde_json = { workspace = true, features = ["std"] } tempfile = { workspace = true } -tokio = { workspace = true, features = ["macros", "rt-multi-thread"] } +tokio = { workspace = true, features = ["io-util", "macros", "rt-multi-thread"] } [[test]] name = "package_boundary" diff --git a/crates/sdk/src/adapters/blossom.rs b/crates/sdk/src/adapters/blossom.rs @@ -0,0 +1,1611 @@ +use std::{net::SocketAddr, time::Duration}; + +use radroots_blossom::{BlobDescriptor, BlobUrl, MediaType}; +use reqwest::{ + StatusCode, + header::{ACCEPT, ACCEPT_ENCODING, AUTHORIZATION, CONTENT_LENGTH, CONTENT_TYPE, LOCATION}, +}; + +use crate::transport::{ + BlossomCancellation, BlossomConfig, BlossomEndpoint, BlossomError, BlossomErrorKind, + BlossomImageDimensions, BlossomPhase, BlossomUploadReceipt, BlossomUploadRequest, +}; + +const MAX_RESOLVED_ADDRESSES: usize = 32; +const X_SHA_256: &str = "x-sha-256"; + +pub(crate) async fn upload( + config: BlossomConfig, + request: BlossomUploadRequest, + authorization: crate::signing::AuthorizationHeader, + cancellation: BlossomCancellation, +) -> Result<BlossomUploadReceipt, BlossomError> { + upload_with_authorization(config, request, authorization.as_str(), cancellation).await +} + +async fn upload_with_authorization( + config: BlossomConfig, + request: BlossomUploadRequest, + authorization: &str, + cancellation: BlossomCancellation, +) -> Result<BlossomUploadReceipt, BlossomError> { + if request.byte_size() > config.max_blob_bytes() { + return Err(failure( + BlossomErrorKind::ResponseTooLarge, + BlossomPhase::Verification, + false, + false, + 0, + )); + } + let endpoint = config + .profile() + .endpoint_for_blob(request.expected_url()) + .cloned() + .ok_or_else(|| { + failure( + BlossomErrorKind::EndpointNotConfigured, + BlossomPhase::Configuration, + false, + false, + 0, + ) + })?; + + let mut upload_attempts = 0_u8; + let descriptor = loop { + ensure_not_cancelled(&cancellation, BlossomPhase::Upload, upload_attempts, false)?; + upload_attempts = upload_attempts.saturating_add(1); + match upload_once( + &config, + &endpoint, + &request, + authorization, + &cancellation, + upload_attempts, + ) + .await + { + Ok(descriptor) => break descriptor, + Err(error) if error.retryable() && upload_attempts < config.max_attempts() => { + retry_delay(&config, upload_attempts, &cancellation, true).await?; + } + Err(error) => return Err(error), + } + }; + + let verified_upload = verify_descriptor(&request, descriptor, upload_attempts)?; + let mut retrieval_attempts = 0_u8; + let retrieved = loop { + ensure_not_cancelled( + &cancellation, + BlossomPhase::Retrieval, + upload_attempts.saturating_add(retrieval_attempts), + true, + )?; + retrieval_attempts = retrieval_attempts.saturating_add(1); + match retrieve_once( + &config, + verified_upload.url().as_blob_url().clone(), + &request, + &cancellation, + upload_attempts.saturating_add(retrieval_attempts), + ) + .await + { + Ok(bytes) => break bytes, + Err(error) if error.retryable() && retrieval_attempts < config.max_attempts() => { + retry_delay(&config, retrieval_attempts, &cancellation, true).await?; + } + Err(error) => return Err(error), + } + }; + + if retrieved.as_slice() != request.bytes() { + return Err(failure( + BlossomErrorKind::RetrievedBytesMismatch, + BlossomPhase::Verification, + false, + true, + upload_attempts.saturating_add(retrieval_attempts), + )); + } + verify_image( + retrieved.as_slice(), + request.media_type(), + request.dimensions(), + ) + .map_err(|error| with_operation(error, true, upload_attempts + retrieval_attempts))?; + let verified_retrieval = verified_upload + .into_descriptor() + .approve_reference() + .and_then(|approved| approved.verify_bytes(retrieved.as_slice(), request.media_type())) + .map_err(|_| { + failure( + BlossomErrorKind::RetrievedBytesMismatch, + BlossomPhase::Verification, + false, + true, + upload_attempts.saturating_add(retrieval_attempts), + ) + })?; + + Ok(BlossomUploadReceipt::new( + verified_retrieval, + request.dimensions(), + upload_attempts.saturating_add(retrieval_attempts), + request.verified_at_unix_ms(), + )) +} + +// Direct DNS/socket/HTTP behavior is verified by the local real-I/O suite; +// deterministic coverage owns the surrounding retry, validation, and durable +// state policy. +#[cfg_attr(coverage_nightly, coverage(off))] +async fn upload_once( + config: &BlossomConfig, + endpoint: &BlossomEndpoint, + request: &BlossomUploadRequest, + authorization: &str, + cancellation: &BlossomCancellation, + attempt: u8, +) -> Result<BlobDescriptor, BlossomError> { + let client = hardened_client( + config, + endpoint, + cancellation, + BlossomPhase::Upload, + attempt, + false, + ) + .await?; + let pending = client + .put(endpoint.upload_url()) + .header(AUTHORIZATION, authorization) + .header(X_SHA_256, request.sha256().to_string()) + .header(CONTENT_TYPE, request.media_type().as_str()) + .header(ACCEPT, "application/json") + .header(ACCEPT_ENCODING, "identity") + .body(request.bytes().to_vec()) + .send(); + let response = tokio::select! { + biased; + _ = cancellation.cancelled() => { + return Err(failure( + BlossomErrorKind::Cancelled, + BlossomPhase::Upload, + true, + true, + attempt, + )); + } + response = pending => response.map_err(|error| request_error(error, BlossomPhase::Upload, true, attempt))?, + }; + + if response.status().is_redirection() { + return Err(failure( + BlossomErrorKind::UnsafeRedirect, + BlossomPhase::Upload, + false, + true, + attempt, + )); + } + if !matches!(response.status(), StatusCode::OK | StatusCode::CREATED) { + return Err(http_status_error( + response.status(), + BlossomPhase::Upload, + true, + attempt, + )); + } + { + let content_type = response + .headers() + .get(CONTENT_TYPE) + .and_then(|value| value.to_str().ok()) + .ok_or_else(|| { + failure( + BlossomErrorKind::InvalidDescriptor, + BlossomPhase::Descriptor, + false, + true, + attempt, + ) + })?; + if content_type + .split(';') + .next() + .is_none_or(|value| !value.trim().eq_ignore_ascii_case("application/json")) + { + return Err(failure( + BlossomErrorKind::InvalidDescriptor, + BlossomPhase::Descriptor, + false, + true, + attempt, + )); + } + } + let bytes = read_bounded( + response, + config.max_descriptor_bytes(), + cancellation, + BlossomPhase::Descriptor, + true, + attempt, + ) + .await?; + serde_json::from_slice(bytes.as_slice()).map_err(|_| { + failure( + BlossomErrorKind::InvalidDescriptor, + BlossomPhase::Descriptor, + false, + true, + attempt, + ) + }) +} + +fn verify_descriptor( + request: &BlossomUploadRequest, + descriptor: BlobDescriptor, + attempts: u8, +) -> Result<radroots_blossom::ByteVerifiedDescriptor, BlossomError> { + // `BlobDescriptor` construction already binds `sha256` to the URL hash, + // so equality of the typed URL proves equality of that hash as well. + if descriptor.url() != request.expected_url() + || descriptor.size() != request.byte_size() + || descriptor.media_type() != request.media_type() + { + return Err(failure( + BlossomErrorKind::DescriptorMismatch, + BlossomPhase::Descriptor, + false, + true, + attempts, + )); + } + descriptor + .approve_reference() + .and_then(|approved| approved.verify_bytes(request.bytes(), request.media_type())) + .map_err(|_| { + failure( + BlossomErrorKind::DescriptorMismatch, + BlossomPhase::Descriptor, + false, + true, + attempts, + ) + }) +} + +#[cfg_attr(coverage_nightly, coverage(off))] +async fn retrieve_once( + config: &BlossomConfig, + mut url: BlobUrl, + request: &BlossomUploadRequest, + cancellation: &BlossomCancellation, + attempt: u8, +) -> Result<Vec<u8>, BlossomError> { + for redirects in 0..=config.max_redirects() { + let endpoint = config.profile().endpoint_for_blob(&url).ok_or_else(|| { + failure( + BlossomErrorKind::UnsafeRedirect, + BlossomPhase::Retrieval, + false, + true, + attempt, + ) + })?; + let client = hardened_client( + config, + endpoint, + cancellation, + BlossomPhase::Retrieval, + attempt, + true, + ) + .await?; + let pending = client + .get(url.as_str()) + .header(ACCEPT, request.media_type().as_str()) + .header(ACCEPT_ENCODING, "identity") + .send(); + let response = tokio::select! { + biased; + _ = cancellation.cancelled() => { + return Err(failure( + BlossomErrorKind::Cancelled, + BlossomPhase::Retrieval, + true, + true, + attempt, + )); + } + response = pending => response.map_err(|error| request_error(error, BlossomPhase::Retrieval, true, attempt))?, + }; + if response.status().is_redirection() { + if redirects == config.max_redirects() { + return Err(failure( + BlossomErrorKind::RedirectLimit, + BlossomPhase::Retrieval, + false, + true, + attempt, + )); + } + let location = response + .headers() + .get(LOCATION) + .and_then(|value| value.to_str().ok()) + .ok_or_else(|| { + failure( + BlossomErrorKind::UnsafeRedirect, + BlossomPhase::Retrieval, + false, + true, + attempt, + ) + })?; + let base = reqwest::Url::parse(url.as_str()).map_err(|_| { + failure( + BlossomErrorKind::UnsafeRedirect, + BlossomPhase::Retrieval, + false, + true, + attempt, + ) + })?; + let next = base.join(location).map_err(|_| { + failure( + BlossomErrorKind::UnsafeRedirect, + BlossomPhase::Retrieval, + false, + true, + attempt, + ) + })?; + let next = BlobUrl::parse(next.as_str()).map_err(|_| { + failure( + BlossomErrorKind::UnsafeRedirect, + BlossomPhase::Retrieval, + false, + true, + attempt, + ) + })?; + if next.hash_path().hash() != request.sha256() + || next.clone().approve().is_err() + || config.profile().endpoint_for_blob(&next).is_none() + { + return Err(failure( + BlossomErrorKind::UnsafeRedirect, + BlossomPhase::Retrieval, + false, + true, + attempt, + )); + } + url = next; + continue; + } + if response.status() != StatusCode::OK { + return Err(http_status_error( + response.status(), + BlossomPhase::Retrieval, + true, + attempt, + )); + } + let actual_media_type = response + .headers() + .get(CONTENT_TYPE) + .and_then(|value| value.to_str().ok()) + .and_then(|value| MediaType::parse(value).ok()) + .ok_or_else(|| { + failure( + BlossomErrorKind::MediaTypeMismatch, + BlossomPhase::Verification, + false, + true, + attempt, + ) + })?; + if &actual_media_type != request.media_type() { + return Err(failure( + BlossomErrorKind::MediaTypeMismatch, + BlossomPhase::Verification, + false, + true, + attempt, + )); + } + if response + .headers() + .get(CONTENT_LENGTH) + .and_then(|value| value.to_str().ok()) + .and_then(|value| value.parse::<u64>().ok()) + .is_some_and(|size| size != request.byte_size() || size > config.max_blob_bytes()) + { + return Err(failure( + BlossomErrorKind::ResponseTooLarge, + BlossomPhase::Retrieval, + false, + true, + attempt, + )); + } + return read_bounded( + response, + usize::try_from(config.max_blob_bytes()).unwrap_or(usize::MAX), + cancellation, + BlossomPhase::Retrieval, + true, + attempt, + ) + .await; + } + Err(failure( + BlossomErrorKind::RedirectLimit, + BlossomPhase::Retrieval, + false, + true, + attempt, + )) +} + +#[cfg_attr(coverage_nightly, coverage(off))] +async fn hardened_client( + config: &BlossomConfig, + endpoint: &BlossomEndpoint, + cancellation: &BlossomCancellation, + phase: BlossomPhase, + attempts: u8, + possible_orphan: bool, +) -> Result<reqwest::Client, BlossomError> { + let addresses = resolve( + endpoint, + config.connect_timeout(), + cancellation, + phase, + attempts, + possible_orphan, + ) + .await?; + let mut builder = reqwest::Client::builder() + .redirect(reqwest::redirect::Policy::none()) + .connect_timeout(config.connect_timeout()) + .timeout(config.request_timeout()) + .pool_max_idle_per_host(0); + if endpoint.host().parse::<std::net::IpAddr>().is_err() { + builder = builder.resolve(endpoint.host(), addresses[0]); + } + builder.build().map_err(|_| { + failure( + BlossomErrorKind::Transport, + phase, + true, + possible_orphan, + attempts, + ) + }) +} + +#[cfg_attr(coverage_nightly, coverage(off))] +async fn resolve( + endpoint: &BlossomEndpoint, + timeout: Duration, + cancellation: &BlossomCancellation, + phase: BlossomPhase, + attempts: u8, + possible_orphan: bool, +) -> Result<Vec<SocketAddr>, BlossomError> { + let lookup = tokio::net::lookup_host((endpoint.host(), endpoint.port())); + let mut resolved = tokio::select! { + biased; + _ = cancellation.cancelled() => { + return Err(failure( + BlossomErrorKind::Cancelled, + phase, + true, + possible_orphan, + attempts, + )); + } + result = tokio::time::timeout(timeout, lookup) => result + .map_err(|_| failure( + BlossomErrorKind::Timeout, + phase, + true, + possible_orphan, + attempts, + ))? + .map_err(|_| { + failure( + BlossomErrorKind::ResolutionFailed, + phase, + true, + possible_orphan, + attempts, + ) + })?, + }; + let addresses = resolved + .by_ref() + .take(MAX_RESOLVED_ADDRESSES + 1) + .collect::<Vec<_>>(); + if addresses.is_empty() || addresses.len() > MAX_RESOLVED_ADDRESSES { + return Err(failure( + BlossomErrorKind::ResolutionFailed, + phase, + true, + possible_orphan, + attempts, + )); + } + endpoint + .validate_resolved_addresses(addresses.iter().map(SocketAddr::ip)) + .map_err(|error| with_operation(error, possible_orphan, attempts))?; + Ok(addresses) +} + +#[cfg_attr(coverage_nightly, coverage(off))] +async fn read_bounded( + mut response: reqwest::Response, + max_bytes: usize, + cancellation: &BlossomCancellation, + phase: BlossomPhase, + possible_orphan: bool, + attempts: u8, +) -> Result<Vec<u8>, BlossomError> { + if response + .content_length() + .is_some_and(|length| length > max_bytes as u64) + { + return Err(failure( + BlossomErrorKind::ResponseTooLarge, + phase, + false, + possible_orphan, + attempts, + )); + } + let mut bytes = Vec::with_capacity( + response + .content_length() + .and_then(|length| usize::try_from(length).ok()) + .unwrap_or(0) + .min(max_bytes), + ); + loop { + let pending = response.chunk(); + let chunk = tokio::select! { + biased; + _ = cancellation.cancelled() => { + return Err(failure( + BlossomErrorKind::Cancelled, + phase, + true, + possible_orphan, + attempts, + )); + } + chunk = pending => chunk.map_err(|error| request_error(error, phase, possible_orphan, attempts))?, + }; + let Some(chunk) = chunk else { + break; + }; + if bytes.len().saturating_add(chunk.len()) > max_bytes { + return Err(failure( + BlossomErrorKind::ResponseTooLarge, + phase, + false, + possible_orphan, + attempts, + )); + } + bytes.extend_from_slice(&chunk); + } + Ok(bytes) +} + +async fn retry_delay( + config: &BlossomConfig, + attempt: u8, + cancellation: &BlossomCancellation, + possible_orphan: bool, +) -> Result<(), BlossomError> { + let exponent = u32::from(attempt.saturating_sub(1)).min(16); + let factor = 1_u32 << exponent; + let delay = config + .initial_retry_delay() + .saturating_mul(factor) + .min(Duration::from_secs(30)); + tokio::select! { + biased; + _ = cancellation.cancelled() => Err(failure( + BlossomErrorKind::Cancelled, + BlossomPhase::Upload, + true, + possible_orphan, + attempt, + )), + _ = tokio::time::sleep(delay) => Ok(()), + } +} + +fn ensure_not_cancelled( + cancellation: &BlossomCancellation, + phase: BlossomPhase, + attempts: u8, + possible_orphan: bool, +) -> Result<(), BlossomError> { + if cancellation.is_cancelled() { + Err(failure( + BlossomErrorKind::Cancelled, + phase, + true, + possible_orphan, + attempts, + )) + } else { + Ok(()) + } +} + +fn request_error( + error: reqwest::Error, + phase: BlossomPhase, + possible_orphan: bool, + attempts: u8, +) -> BlossomError { + let kind = if error.is_timeout() { + BlossomErrorKind::Timeout + } else { + BlossomErrorKind::Transport + }; + failure(kind, phase, true, possible_orphan, attempts) +} + +fn http_status_error( + status: StatusCode, + phase: BlossomPhase, + possible_orphan: bool, + attempts: u8, +) -> BlossomError { + let retryable = matches!( + status, + StatusCode::REQUEST_TIMEOUT + | StatusCode::TOO_EARLY + | StatusCode::TOO_MANY_REQUESTS + | StatusCode::INTERNAL_SERVER_ERROR + | StatusCode::BAD_GATEWAY + | StatusCode::SERVICE_UNAVAILABLE + | StatusCode::GATEWAY_TIMEOUT + ); + failure( + BlossomErrorKind::HttpStatus, + phase, + retryable, + possible_orphan, + attempts, + ) +} + +fn with_operation(error: BlossomError, possible_orphan: bool, attempts: u8) -> BlossomError { + error.with_operation(possible_orphan, attempts) +} + +fn failure( + kind: BlossomErrorKind, + phase: BlossomPhase, + retryable: bool, + possible_orphan: bool, + attempts: u8, +) -> BlossomError { + BlossomError::new(kind, phase, retryable, possible_orphan, attempts) +} + +pub(crate) fn verify_image( + bytes: &[u8], + media_type: &MediaType, + expected: BlossomImageDimensions, +) -> Result<(), BlossomError> { + let detected = detect_image(bytes).ok_or_else(|| { + failure( + BlossomErrorKind::InvalidImageBytes, + BlossomPhase::Verification, + false, + false, + 0, + ) + })?; + let declared = match media_type.as_str() { + "image/png" => ImageKind::Png, + "image/jpeg" => ImageKind::Jpeg, + "image/gif" => ImageKind::Gif, + "image/webp" => ImageKind::Webp, + _ => { + return Err(failure( + BlossomErrorKind::UnsupportedMediaType, + BlossomPhase::Verification, + false, + false, + 0, + )); + } + }; + if declared != detected.0 { + return Err(failure( + BlossomErrorKind::MediaTypeMismatch, + BlossomPhase::Verification, + false, + false, + 0, + )); + } + if detected.1 != expected { + return Err(failure( + BlossomErrorKind::DimensionMismatch, + BlossomPhase::Verification, + false, + false, + 0, + )); + } + Ok(()) +} + +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +enum ImageKind { + Png, + Jpeg, + Gif, + Webp, +} + +fn detect_image(bytes: &[u8]) -> Option<(ImageKind, BlossomImageDimensions)> { + detect_png(bytes) + .map(|dimensions| (ImageKind::Png, dimensions)) + .or_else(|| detect_jpeg(bytes).map(|dimensions| (ImageKind::Jpeg, dimensions))) + .or_else(|| detect_gif(bytes).map(|dimensions| (ImageKind::Gif, dimensions))) + .or_else(|| detect_webp(bytes).map(|dimensions| (ImageKind::Webp, dimensions))) +} + +fn detect_png(bytes: &[u8]) -> Option<BlossomImageDimensions> { + if bytes.len() < 24 || &bytes[..8] != b"\x89PNG\r\n\x1a\n" || &bytes[12..16] != b"IHDR" { + return None; + } + BlossomImageDimensions::new( + u32::from_be_bytes(bytes[16..20].try_into().ok()?), + u32::from_be_bytes(bytes[20..24].try_into().ok()?), + ) + .ok() +} + +fn detect_gif(bytes: &[u8]) -> Option<BlossomImageDimensions> { + if bytes.len() < 10 || !matches!(&bytes[..6], b"GIF87a" | b"GIF89a") { + return None; + } + BlossomImageDimensions::new( + u32::from(u16::from_le_bytes(bytes[6..8].try_into().ok()?)), + u32::from(u16::from_le_bytes(bytes[8..10].try_into().ok()?)), + ) + .ok() +} + +fn detect_jpeg(bytes: &[u8]) -> Option<BlossomImageDimensions> { + if bytes.len() < 4 || bytes[..2] != [0xff, 0xd8] { + return None; + } + let mut offset = 2_usize; + while offset < bytes.len() { + while bytes.get(offset) == Some(&0xff) { + offset += 1; + } + let marker = *bytes.get(offset)?; + offset += 1; + if marker == 0xd9 || marker == 0xda { + return None; + } + if marker == 0x01 || (0xd0..=0xd7).contains(&marker) { + continue; + } + let length = usize::from(u16::from_be_bytes([ + *bytes.get(offset)?, + *bytes.get(offset + 1)?, + ])); + if length < 2 || offset.checked_add(length)? > bytes.len() { + return None; + } + if matches!( + marker, + 0xc0 | 0xc1 + | 0xc2 + | 0xc3 + | 0xc5 + | 0xc6 + | 0xc7 + | 0xc9 + | 0xca + | 0xcb + | 0xcd + | 0xce + | 0xcf + ) { + if length < 7 { + return None; + } + return BlossomImageDimensions::new( + u32::from(u16::from_be_bytes([ + *bytes.get(offset + 5)?, + *bytes.get(offset + 6)?, + ])), + u32::from(u16::from_be_bytes([ + *bytes.get(offset + 3)?, + *bytes.get(offset + 4)?, + ])), + ) + .ok(); + } + offset += length; + } + None +} + +fn detect_webp(bytes: &[u8]) -> Option<BlossomImageDimensions> { + if bytes.len() < 30 || &bytes[..4] != b"RIFF" || &bytes[8..12] != b"WEBP" { + return None; + } + match &bytes[12..16] { + b"VP8X" => BlossomImageDimensions::new( + little_u24(&bytes[24..27])?.checked_add(1)?, + little_u24(&bytes[27..30])?.checked_add(1)?, + ) + .ok(), + b"VP8L" if bytes.get(20) == Some(&0x2f) => { + let packed = u32::from_le_bytes(bytes[21..25].try_into().ok()?); + BlossomImageDimensions::new( + (packed & 0x3fff).checked_add(1)?, + ((packed >> 14) & 0x3fff).checked_add(1)?, + ) + .ok() + } + b"VP8 " if bytes.get(23..26) == Some(&[0x9d, 0x01, 0x2a]) => BlossomImageDimensions::new( + u32::from(u16::from_le_bytes(bytes[26..28].try_into().ok()?) & 0x3fff), + u32::from(u16::from_le_bytes(bytes[28..30].try_into().ok()?) & 0x3fff), + ) + .ok(), + _ => None, + } +} + +fn little_u24(bytes: &[u8]) -> Option<u32> { + Some( + u32::from(*bytes.first()?) + | u32::from(*bytes.get(1)?) << 8 + | u32::from(*bytes.get(2)?) << 16, + ) +} + +#[cfg(test)] +mod tests { + use super::*; + use radroots_blossom::Sha256; + use std::sync::Arc; + use tokio::io::{AsyncReadExt, AsyncWriteExt}; + use tokio::net::TcpListener; + + fn png(width: u32, height: u32) -> Vec<u8> { + let mut bytes = b"\x89PNG\r\n\x1a\n\0\0\0\rIHDR".to_vec(); + bytes.extend_from_slice(&width.to_be_bytes()); + bytes.extend_from_slice(&height.to_be_bytes()); + bytes + } + + #[test] + fn image_headers_bind_mime_and_dimensions() { + let dimensions = BlossomImageDimensions::new(1200, 900).expect("dimensions"); + let bytes = png(1200, 900); + assert!(verify_image(&bytes, &MediaType::parse("image/png").unwrap(), dimensions).is_ok()); + assert_eq!( + verify_image(&bytes, &MediaType::parse("image/jpeg").unwrap(), dimensions) + .expect_err("wrong MIME") + .kind(), + BlossomErrorKind::MediaTypeMismatch + ); + assert_eq!( + verify_image( + &bytes, + &MediaType::parse("image/png").unwrap(), + BlossomImageDimensions::new(1, 1).unwrap(), + ) + .expect_err("wrong dimensions") + .kind(), + BlossomErrorKind::DimensionMismatch + ); + assert_eq!( + verify_image( + b"not an image", + &MediaType::parse("image/png").unwrap(), + dimensions + ) + .expect_err("invalid image") + .kind(), + BlossomErrorKind::InvalidImageBytes + ); + } + + #[test] + fn supported_image_headers_are_bounded_and_nonzero() { + let gif = b"GIF89a\x02\0\x03\0"; + assert_eq!( + detect_gif(gif), + Some(BlossomImageDimensions::new(2, 3).unwrap()) + ); + + let jpeg = [ + 0xff, 0xd8, 0xff, 0xc0, 0x00, 0x11, 0x08, 0x00, 0x03, 0x00, 0x02, 0x03, 0x01, 0x11, + 0x00, 0x02, 0x11, 0x00, 0x03, 0x11, 0x00, + ]; + assert_eq!( + detect_jpeg(&jpeg), + Some(BlossomImageDimensions::new(2, 3).unwrap()) + ); + + let mut webp = vec![0_u8; 30]; + webp[..4].copy_from_slice(b"RIFF"); + webp[8..12].copy_from_slice(b"WEBP"); + webp[12..16].copy_from_slice(b"VP8X"); + webp[24..27].copy_from_slice(&[1, 0, 0]); + webp[27..30].copy_from_slice(&[2, 0, 0]); + assert_eq!( + detect_webp(&webp), + Some(BlossomImageDimensions::new(2, 3).unwrap()) + ); + } + + #[test] + fn image_header_decoders_reject_every_malformed_boundary() { + let mut invalid_png_signature = vec![0_u8; 24]; + invalid_png_signature[12..16].copy_from_slice(b"IHDR"); + let mut invalid_png_chunk = png(2, 3); + invalid_png_chunk[12..16].copy_from_slice(b"NOPE"); + let zero_png = png(0, 3); + assert_eq!(detect_png(b"short"), None); + assert_eq!(detect_png(&invalid_png_signature), None); + assert_eq!(detect_png(&invalid_png_chunk), None); + assert_eq!(detect_png(&zero_png), None); + + assert_eq!(detect_gif(b"short"), None); + assert_eq!(detect_gif(b"GIF00a\x02\0\x03\0"), None); + assert_eq!( + detect_gif(b"GIF87a\x02\0\x03\0"), + Some(BlossomImageDimensions::new(2, 3).unwrap()) + ); + assert_eq!(detect_gif(b"GIF89a\0\0\x03\0"), None); + + for malformed in [ + vec![0xff, 0xd8], + vec![0, 0, 0, 0], + vec![0xff, 0xd8, 0xff, 0xd9], + vec![0xff, 0xd8, 0xff, 0xda], + vec![0xff, 0xd8, 0x01], + vec![0xff, 0xd8, 0xd0], + vec![0xff, 0xd8, 0xff, 0xe0, 0, 1], + vec![0xff, 0xd8, 0xff, 0xe0, 0, 100], + vec![0xff, 0xd8, 0xff, 0xc0, 0, 6, 0, 0, 0, 0], + ] { + assert_eq!(detect_jpeg(&malformed), None, "{malformed:?}"); + } + let jpeg_with_prefix = [ + 0xff, 0xd8, 0xff, 0xe0, 0, 2, 0xff, 0xc2, 0, 7, 8, 0, 3, 0, 2, + ]; + assert_eq!( + detect_jpeg(&jpeg_with_prefix), + Some(BlossomImageDimensions::new(2, 3).unwrap()) + ); + + assert_eq!(detect_webp(b"short"), None); + let mut wrong_riff = vec![0_u8; 30]; + wrong_riff[8..12].copy_from_slice(b"WEBP"); + assert_eq!(detect_webp(&wrong_riff), None); + let mut wrong_webp = vec![0_u8; 30]; + wrong_webp[..4].copy_from_slice(b"RIFF"); + assert_eq!(detect_webp(&wrong_webp), None); + + let mut lossless = vec![0_u8; 30]; + lossless[..4].copy_from_slice(b"RIFF"); + lossless[8..12].copy_from_slice(b"WEBP"); + lossless[12..16].copy_from_slice(b"VP8L"); + lossless[20] = 0x2f; + let packed = 1_u32 | (2_u32 << 14); + lossless[21..25].copy_from_slice(&packed.to_le_bytes()); + assert_eq!( + detect_webp(&lossless), + Some(BlossomImageDimensions::new(2, 3).unwrap()) + ); + lossless[20] = 0; + assert_eq!(detect_webp(&lossless), None); + + let mut lossy = vec![0_u8; 30]; + lossy[..4].copy_from_slice(b"RIFF"); + lossy[8..12].copy_from_slice(b"WEBP"); + lossy[12..16].copy_from_slice(b"VP8 "); + lossy[23..26].copy_from_slice(&[0x9d, 0x01, 0x2a]); + lossy[26..28].copy_from_slice(&2_u16.to_le_bytes()); + lossy[28..30].copy_from_slice(&3_u16.to_le_bytes()); + assert_eq!( + detect_webp(&lossy), + Some(BlossomImageDimensions::new(2, 3).unwrap()) + ); + lossy[23] = 0; + assert_eq!(detect_webp(&lossy), None); + lossy[12..16].copy_from_slice(b"NOPE"); + assert_eq!(detect_webp(&lossy), None); + + assert_eq!(little_u24(&[]), None); + assert_eq!(little_u24(&[1]), None); + assert_eq!(little_u24(&[1, 2]), None); + assert_eq!(little_u24(&[1, 2, 3]), Some(0x03_02_01)); + } + + #[test] + fn image_verification_supports_every_declared_mime() { + let gif = b"GIF89a\x02\0\x03\0"; + assert!( + verify_image( + gif, + &MediaType::parse("image/gif").unwrap(), + BlossomImageDimensions::new(2, 3).unwrap(), + ) + .is_ok() + ); + + let jpeg = [0xff, 0xd8, 0xff, 0xc0, 0, 7, 8, 0, 3, 0, 2]; + assert!( + verify_image( + &jpeg, + &MediaType::parse("image/jpeg").unwrap(), + BlossomImageDimensions::new(2, 3).unwrap(), + ) + .is_ok() + ); + + let mut webp = vec![0_u8; 30]; + webp[..4].copy_from_slice(b"RIFF"); + webp[8..12].copy_from_slice(b"WEBP"); + webp[12..16].copy_from_slice(b"VP8X"); + webp[24..27].copy_from_slice(&[1, 0, 0]); + webp[27..30].copy_from_slice(&[2, 0, 0]); + assert!( + verify_image( + &webp, + &MediaType::parse("image/webp").unwrap(), + BlossomImageDimensions::new(2, 3).unwrap(), + ) + .is_ok() + ); + assert_eq!( + verify_image( + &webp, + &MediaType::parse("application/octet-stream").unwrap(), + BlossomImageDimensions::new(2, 3).unwrap(), + ) + .expect_err("unsupported media type") + .kind(), + BlossomErrorKind::UnsupportedMediaType + ); + } + + #[tokio::test] + async fn retry_and_error_classification_are_bounded_and_redacted() { + for status in [ + StatusCode::REQUEST_TIMEOUT, + StatusCode::TOO_EARLY, + StatusCode::TOO_MANY_REQUESTS, + StatusCode::INTERNAL_SERVER_ERROR, + StatusCode::BAD_GATEWAY, + StatusCode::SERVICE_UNAVAILABLE, + StatusCode::GATEWAY_TIMEOUT, + ] { + assert!(http_status_error(status, BlossomPhase::Upload, true, 1).retryable()); + } + assert!( + !http_status_error(StatusCode::BAD_REQUEST, BlossomPhase::Upload, false, 1).retryable() + ); + + let cancellation = BlossomCancellation::default(); + assert!(ensure_not_cancelled(&cancellation, BlossomPhase::Upload, 0, false).is_ok()); + cancellation.cancel(); + assert_eq!( + ensure_not_cancelled(&cancellation, BlossomPhase::Retrieval, 2, true) + .expect_err("cancelled") + .kind(), + BlossomErrorKind::Cancelled + ); + + let profile = + crate::transport::BlossomProfile::simulator(["http://127.0.0.1:9"]).expect("profile"); + let config = BlossomConfig::from_profile(profile) + .with_network_policy( + Duration::from_millis(10), + Duration::from_millis(10), + 1, + Duration::from_millis(1), + ) + .unwrap(); + let delay = BlossomCancellation::default(); + assert!(retry_delay(&config, 1, &delay, false).await.is_ok()); + delay.cancel(); + assert_eq!( + retry_delay(&config, 20, &delay, true) + .await + .expect_err("cancelled delay") + .kind(), + BlossomErrorKind::Cancelled + ); + + let transport_error = reqwest::Client::new() + .get("http://127.0.0.1:9") + .send() + .await + .expect_err("closed port"); + assert_eq!( + request_error(transport_error, BlossomPhase::Upload, false, 1).kind(), + BlossomErrorKind::Transport + ); + + let bytes = png(2, 3); + let request = upload_request("http://127.0.0.1:9", bytes.clone()); + let too_small = BlossomConfig::from_profile( + crate::transport::BlossomProfile::simulator(["http://127.0.0.1:9"]).unwrap(), + ) + .with_limits(1, 100, 0) + .unwrap(); + assert_eq!( + upload_with_authorization( + too_small, + request, + "Nostr redacted", + BlossomCancellation::default(), + ) + .await + .expect_err("oversize request") + .kind(), + BlossomErrorKind::ResponseTooLarge + ); + let unconfigured_request = upload_request("http://127.0.0.1:10", bytes); + assert_eq!( + upload_with_authorization( + config, + unconfigured_request, + "Nostr redacted", + BlossomCancellation::default(), + ) + .await + .expect_err("unconfigured origin") + .kind(), + BlossomErrorKind::EndpointNotConfigured + ); + } + + #[test] + fn descriptor_verification_checks_each_identity_field() { + let bytes = png(2, 3); + let request = upload_request("http://127.0.0.1:3000", bytes); + let descriptor = |url: BlobUrl, size: u64, media_type: &str| { + BlobDescriptor::new( + url, + request.sha256(), + size, + MediaType::parse(media_type).unwrap(), + 1, + ) + .unwrap() + }; + assert!( + verify_descriptor( + &request, + descriptor( + BlobUrl::parse( + format!("http://localhost:3000/{}.png", request.sha256()).as_str() + ) + .unwrap(), + request.byte_size(), + "image/png", + ), + 1, + ) + .is_err() + ); + assert!( + verify_descriptor( + &request, + descriptor( + request.expected_url().clone(), + request.byte_size() + 1, + "image/png" + ), + 1, + ) + .is_err() + ); + assert!( + verify_descriptor( + &request, + descriptor( + request.expected_url().clone(), + request.byte_size(), + "image/jpeg", + ), + 1, + ) + .is_err() + ); + } + + #[derive(Clone, Copy)] + enum RetrievalResponse { + Exact, + Altered, + RedirectExternal, + WrongMime, + Oversize, + Stall, + } + + async fn read_request(stream: &mut tokio::net::TcpStream) -> Vec<u8> { + let mut request = Vec::new(); + let header_end = loop { + let mut chunk = [0_u8; 1024]; + let read = stream.read(&mut chunk).await.expect("request read"); + assert_ne!(read, 0, "request ended before headers"); + request.extend_from_slice(&chunk[..read]); + if let Some(end) = request.windows(4).position(|value| value == b"\r\n\r\n") { + break end + 4; + } + assert!(request.len() < 64 * 1024, "request headers are bounded"); + }; + let headers = String::from_utf8_lossy(&request[..header_end]); + let content_length = headers + .lines() + .find_map(|line| { + line.to_ascii_lowercase() + .strip_prefix("content-length: ") + .and_then(|value| value.trim().parse::<usize>().ok()) + }) + .unwrap_or(0); + while request.len() - header_end < content_length { + let mut chunk = [0_u8; 1024]; + let read = stream.read(&mut chunk).await.expect("body read"); + assert_ne!(read, 0, "request ended before body"); + request.extend_from_slice(&chunk[..read]); + } + request + } + + async fn spawn_server( + bytes: Vec<u8>, + retrieval: RetrievalResponse, + bad_descriptor: bool, + ) -> (String, tokio::task::JoinHandle<Vec<u8>>) { + let listener = TcpListener::bind("127.0.0.1:0").await.expect("bind"); + let address = listener.local_addr().expect("address"); + let origin = format!("http://{address}"); + let expected_hash = Sha256::digest(bytes.as_slice()); + let descriptor_hash = if bad_descriptor { + Sha256::digest(b"different") + } else { + expected_hash + }; + let descriptor_url = format!("{origin}/{descriptor_hash}.png"); + let descriptor = BlobDescriptor::new( + BlobUrl::parse(descriptor_url.as_str()).expect("url"), + descriptor_hash, + bytes.len() as u64, + MediaType::parse("image/png").expect("media type"), + 1_900_000_000, + ) + .expect("descriptor"); + let descriptor_json = serde_json::to_vec(&descriptor).expect("descriptor json"); + let task = tokio::spawn(async move { + let (mut upload, _) = listener.accept().await.expect("upload accept"); + let upload_request = read_request(&mut upload).await; + let response = format!( + "HTTP/1.1 200 OK\r\nContent-Type: application/json\r\nContent-Length: {}\r\nConnection: close\r\n\r\n", + descriptor_json.len() + ); + upload + .write_all(response.as_bytes()) + .await + .expect("upload head"); + upload + .write_all(&descriptor_json) + .await + .expect("upload body"); + upload.shutdown().await.expect("upload close"); + + if bad_descriptor { + return upload_request; + } + let (mut retrieval_stream, _) = listener.accept().await.expect("retrieval accept"); + let _ = read_request(&mut retrieval_stream).await; + match retrieval { + RetrievalResponse::Exact | RetrievalResponse::Altered => { + let body = if matches!(retrieval, RetrievalResponse::Altered) { + let mut altered = bytes.clone(); + altered[0] ^= 1; + altered + } else { + bytes + }; + let response = format!( + "HTTP/1.1 200 OK\r\nContent-Type: image/png\r\nContent-Length: {}\r\nConnection: close\r\n\r\n", + body.len() + ); + retrieval_stream + .write_all(response.as_bytes()) + .await + .expect("retrieval head"); + retrieval_stream + .write_all(&body) + .await + .expect("retrieval body"); + } + RetrievalResponse::RedirectExternal => { + let location = format!("https://example.com/{expected_hash}.png"); + let response = format!( + "HTTP/1.1 302 Found\r\nLocation: {location}\r\nContent-Length: 0\r\nConnection: close\r\n\r\n" + ); + retrieval_stream + .write_all(response.as_bytes()) + .await + .expect("redirect"); + } + RetrievalResponse::WrongMime => { + let response = format!( + "HTTP/1.1 200 OK\r\nContent-Type: image/jpeg\r\nContent-Length: {}\r\nConnection: close\r\n\r\n", + bytes.len() + ); + retrieval_stream + .write_all(response.as_bytes()) + .await + .expect("wrong MIME"); + retrieval_stream + .write_all(&bytes) + .await + .expect("retrieval body"); + } + RetrievalResponse::Oversize => { + retrieval_stream + .write_all(b"HTTP/1.1 200 OK\r\nContent-Type: image/png\r\nContent-Length: 999999\r\nConnection: close\r\n\r\n") + .await + .expect("oversize"); + } + RetrievalResponse::Stall => { + tokio::time::sleep(Duration::from_secs(1)).await; + } + } + retrieval_stream.shutdown().await.expect("retrieval close"); + upload_request + }); + (origin, task) + } + + async fn spawn_retry_server(bytes: Vec<u8>) -> (String, tokio::task::JoinHandle<()>) { + let listener = TcpListener::bind("127.0.0.1:0").await.expect("bind"); + let address = listener.local_addr().expect("address"); + let origin = format!("http://{address}"); + let hash = Sha256::digest(bytes.as_slice()); + let descriptor = BlobDescriptor::new( + BlobUrl::parse(format!("{origin}/{hash}.png").as_str()).expect("url"), + hash, + bytes.len() as u64, + MediaType::parse("image/png").expect("media type"), + 1_900_000_000, + ) + .expect("descriptor"); + let descriptor_json = serde_json::to_vec(&descriptor).expect("descriptor JSON"); + let task = tokio::spawn(async move { + for step in 0..4 { + let (mut stream, _) = listener.accept().await.expect("accept"); + let _ = read_request(&mut stream).await; + match step { + 0 => { + stream + .write_all(b"HTTP/1.1 503 Service Unavailable\r\nContent-Length: 0\r\nConnection: close\r\n\r\n") + .await + .expect("upload retry"); + } + 1 => { + let head = format!( + "HTTP/1.1 200 OK\r\nContent-Type: application/json\r\nContent-Length: {}\r\nConnection: close\r\n\r\n", + descriptor_json.len() + ); + stream + .write_all(head.as_bytes()) + .await + .expect("descriptor head"); + stream + .write_all(&descriptor_json) + .await + .expect("descriptor body"); + } + 2 => { + stream + .write_all(b"HTTP/1.1 503 Service Unavailable\r\nContent-Length: 0\r\nConnection: close\r\n\r\n") + .await + .expect("retrieval retry"); + } + 3 => { + let head = format!( + "HTTP/1.1 200 OK\r\nContent-Type: image/png\r\nContent-Length: {}\r\nConnection: close\r\n\r\n", + bytes.len() + ); + stream + .write_all(head.as_bytes()) + .await + .expect("retrieval head"); + stream.write_all(&bytes).await.expect("retrieval body"); + } + _ => unreachable!(), + } + stream.shutdown().await.expect("close"); + } + }); + (origin, task) + } + + fn upload_request(origin: &str, bytes: Vec<u8>) -> BlossomUploadRequest { + let hash = Sha256::digest(bytes.as_slice()); + BlossomUploadRequest::new( + BlobUrl::parse(format!("{origin}/{hash}.png").as_str()).expect("blob URL"), + Arc::from(bytes), + MediaType::parse("image/png").expect("media type"), + BlossomImageDimensions::new(2, 3).expect("dimensions"), + 1_900_000_000_000, + ) + .expect("upload request") + } + + fn config(origin: &str) -> BlossomConfig { + BlossomConfig::from_profile( + crate::transport::BlossomProfile::simulator([origin]).expect("profile"), + ) + .with_network_policy( + Duration::from_millis(100), + Duration::from_millis(100), + 1, + Duration::from_millis(1), + ) + .expect("network policy") + } + + #[tokio::test] + async fn loopback_upload_preserves_exact_bytes_and_verifies_retrieval() { + let bytes = png(2, 3); + let (origin, server) = spawn_server(bytes.clone(), RetrievalResponse::Exact, false).await; + let receipt = upload_with_authorization( + config(origin.as_str()), + upload_request(origin.as_str(), bytes.clone()), + "Nostr secret-token-value", + BlossomCancellation::default(), + ) + .await + .expect("verified upload"); + assert_eq!(receipt.descriptor().size(), bytes.len() as u64); + assert_eq!( + receipt.dimensions(), + BlossomImageDimensions::new(2, 3).unwrap() + ); + let upload_wire = server.await.expect("server"); + let header_end = upload_wire + .windows(4) + .position(|value| value == b"\r\n\r\n") + .unwrap() + + 4; + assert_eq!(&upload_wire[header_end..], bytes.as_slice()); + let headers = String::from_utf8_lossy(&upload_wire[..header_end]); + assert!(headers.contains("authorization: Nostr secret-token-value")); + assert!(!format!("{receipt:?}").contains("secret-token-value")); + } + + #[tokio::test] + async fn retryable_upload_and_retrieval_failures_recover_with_bounded_attempts() { + let bytes = png(2, 3); + let (origin, server) = spawn_retry_server(bytes.clone()).await; + let config = BlossomConfig::from_profile( + crate::transport::BlossomProfile::simulator([origin.as_str()]).expect("profile"), + ) + .with_network_policy( + Duration::from_millis(100), + Duration::from_millis(100), + 2, + Duration::from_millis(1), + ) + .expect("network policy"); + let receipt = upload_with_authorization( + config, + upload_request(origin.as_str(), bytes), + "Nostr redacted", + BlossomCancellation::default(), + ) + .await + .expect("retry recovery"); + assert_eq!(receipt.attempts(), 4); + server.await.expect("server"); + } + + #[tokio::test] + async fn descriptor_redirect_body_mime_timeout_and_cancellation_fail_closed() { + let cases = [ + ( + RetrievalResponse::Exact, + true, + BlossomErrorKind::DescriptorMismatch, + ), + ( + RetrievalResponse::RedirectExternal, + false, + BlossomErrorKind::UnsafeRedirect, + ), + ( + RetrievalResponse::Altered, + false, + BlossomErrorKind::RetrievedBytesMismatch, + ), + ( + RetrievalResponse::WrongMime, + false, + BlossomErrorKind::MediaTypeMismatch, + ), + ( + RetrievalResponse::Oversize, + false, + BlossomErrorKind::ResponseTooLarge, + ), + (RetrievalResponse::Stall, false, BlossomErrorKind::Timeout), + ]; + for (response, bad_descriptor, expected) in cases { + let bytes = png(2, 3); + let (origin, server) = spawn_server(bytes.clone(), response, bad_descriptor).await; + let error = upload_with_authorization( + config(origin.as_str()), + upload_request(origin.as_str(), bytes), + "Nostr redacted", + BlossomCancellation::default(), + ) + .await + .expect_err("must fail closed"); + assert_eq!(error.kind(), expected); + assert!(error.possible_orphan()); + assert!(!format!("{error:?}").contains("Nostr redacted")); + server.await.expect("server"); + } + + let cancellation = BlossomCancellation::default(); + cancellation.cancel(); + let bytes = png(2, 3); + let hash = Sha256::digest(bytes.as_slice()); + let request = BlossomUploadRequest::new( + BlobUrl::parse(format!("http://127.0.0.1:9/{hash}.png").as_str()).unwrap(), + Arc::from(bytes), + MediaType::parse("image/png").unwrap(), + BlossomImageDimensions::new(2, 3).unwrap(), + 1, + ) + .unwrap(); + let error = upload_with_authorization( + config("http://127.0.0.1:9"), + request, + "Nostr redacted", + cancellation, + ) + .await + .expect_err("cancelled"); + assert_eq!(error.kind(), BlossomErrorKind::Cancelled); + assert!(!error.possible_orphan()); + } +} diff --git a/crates/sdk/src/adapters/mod.rs b/crates/sdk/src/adapters/mod.rs @@ -1,2 +1,4 @@ +#[cfg(feature = "blossom")] +pub(crate) mod blossom; #[cfg(feature = "radrootsd")] pub(crate) mod radrootsd; diff --git a/crates/sdk/src/client.rs b/crates/sdk/src/client.rs @@ -40,6 +40,8 @@ pub struct ClientBuilder { host_sync: Option<crate::sync::HostPolicy>, #[cfg(feature = "nostr")] nostr: Option<crate::transport::NostrSlot>, + #[cfg(feature = "blossom")] + blossom: Option<crate::transport::BlossomSlot>, capability_availability: BTreeMap<CapabilityId, Availability>, explicitly_configured_capabilities: BTreeSet<CapabilityId>, } @@ -53,6 +55,8 @@ struct ClientInner { sync: Option<radroots_sync::Engine>, #[cfg(feature = "nostr")] nostr: Option<crate::transport::NostrSlot>, + #[cfg(feature = "blossom")] + blossom: Option<crate::transport::BlossomSlot>, capability_availability: BTreeMap<CapabilityId, Availability>, explicitly_configured_capabilities: BTreeSet<CapabilityId>, lifecycle: AtomicU8, @@ -185,6 +189,14 @@ impl ClientBuilder { self } + /// Installs one host-reconfigurable Blossom HTTP adapter slot. + #[cfg(feature = "blossom")] + #[must_use] + pub fn blossom(mut self, slot: crate::transport::BlossomSlot) -> Self { + self.blossom = Some(slot); + self + } + /// Injects an explicitly composed synchronization engine. #[cfg(feature = "sync")] #[must_use] @@ -249,6 +261,8 @@ impl ClientBuilder { sync: self.sync, #[cfg(feature = "nostr")] nostr: self.nostr, + #[cfg(feature = "blossom")] + blossom: self.blossom, capability_availability: self.capability_availability, explicitly_configured_capabilities: self.explicitly_configured_capabilities, lifecycle: AtomicU8::new(OPEN), @@ -387,6 +401,25 @@ impl Client { .configure(profile) } + /// Returns the explicitly composed Blossom adapter, when configured. + #[cfg(feature = "blossom")] + pub fn blossom(&self) -> Result<Option<&crate::transport::BlossomSlot>> { + self.require_open()?; + Ok(self.inner.blossom.as_ref()) + } + + /// Atomically replaces the active Blossom profile without network I/O. + #[cfg(feature = "blossom")] + pub fn configure_blossom(&self, config: crate::transport::BlossomConfig) -> Result<()> { + self.require_open()?; + self.inner + .blossom + .as_ref() + .ok_or_else(Error::shared_operation_unavailable)? + .configure(config) + .map_err(Error::invalid_host_configuration) + } + /// Returns client-scoped canonical synchronization operations, when configured. #[cfg(feature = "sync")] pub fn sync(&self) -> Result<Option<crate::sync::Operations<'_>>> { @@ -516,11 +549,14 @@ impl Drop for CloseAttempt { impl std::fmt::Debug for Client { fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { - formatter - .debug_struct("Client") + let mut debug = formatter.debug_struct("Client"); + debug .field("signer", &self.inner.signer.is_some()) .field("source", &self.inner.source.is_some()) - .field("sink", &self.inner.sink.is_some()) + .field("sink", &self.inner.sink.is_some()); + #[cfg(feature = "blossom")] + debug.field("blossom", &self.inner.blossom.is_some()); + debug .field("closed", &self.is_closed()) .finish_non_exhaustive() } @@ -528,13 +564,15 @@ impl std::fmt::Debug for Client { impl std::fmt::Debug for ClientBuilder { fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { - formatter - .debug_struct("ClientBuilder") + let mut debug = formatter.debug_struct("ClientBuilder"); + debug .field("storage", &self.storage.is_some()) .field("signer", &self.signer.is_some()) .field("source", &self.source.is_some()) - .field("sink", &self.sink.is_some()) - .finish_non_exhaustive() + .field("sink", &self.sink.is_some()); + #[cfg(feature = "blossom")] + debug.field("blossom", &self.blossom.is_some()); + debug.finish_non_exhaustive() } } @@ -697,10 +735,12 @@ mod tests { .build() .expect("outbound client"); assert!(client.signer().expect("signer capability").is_some()); - assert_eq!( - format!("{client:?}"), + let expected = if cfg!(feature = "blossom") { + "Client { signer: true, source: false, sink: true, blossom: false, closed: false, .. }" + } else { "Client { signer: true, source: false, sink: true, closed: false, .. }" - ); + }; + assert_eq!(format!("{client:?}"), expected); } #[cfg(all(feature = "sync", feature = "nip46"))] diff --git a/crates/sdk/src/error.rs b/crates/sdk/src/error.rs @@ -283,7 +283,7 @@ impl Error { } } - #[cfg(any(feature = "sync", feature = "nostr"))] + #[cfg(any(feature = "blossom", feature = "sync", feature = "nostr"))] pub(crate) fn invalid_host_configuration( source: impl error::Error + Send + Sync + 'static, ) -> Self { @@ -298,7 +298,7 @@ impl Error { Self::without_source(ErrorKind::InvalidHostConfiguration) } - #[cfg(any(feature = "sync", feature = "nostr"))] + #[cfg(any(feature = "blossom", feature = "sync", feature = "nostr"))] pub(crate) fn shared_operation_unavailable() -> Self { Self::without_source(ErrorKind::SharedOperationUnavailable) } diff --git a/crates/sdk/src/lib.rs b/crates/sdk/src/lib.rs @@ -25,6 +25,7 @@ #![forbid(unsafe_code)] #![doc = include_str!("../README.md")] +#![cfg_attr(coverage_nightly, feature(coverage_attribute))] mod adapters; diff --git a/crates/sdk/src/transport.rs b/crates/sdk/src/transport.rs @@ -9,9 +9,23 @@ use radroots_transport::{ capability::{Availability, Maturity, SinkCapabilities, SourceCapabilities}, policy::SatisfactionPolicy, }; -#[cfg(feature = "nostr")] +#[cfg(any(feature = "blossom", feature = "nostr"))] use std::sync::{Arc, RwLock}; +#[cfg(feature = "blossom")] +use std::{ + collections::BTreeSet, + net::{IpAddr, Ipv4Addr, Ipv6Addr}, + sync::atomic::{AtomicBool, Ordering}, + time::Duration, +}; + +#[cfg(feature = "blossom")] +use radroots_blossom::{ + BlobUrl, ByteVerifiedDescriptor, MediaType, Sha256, + authorization::{AuthoredUploadClaim, AuthorizationContent, ServerDomain}, +}; + #[cfg(feature = "nostr")] pub use radroots_transport_nostr::{ DEFAULT_PUBLIC_RELAY, ReconnectBackoff, RelayAccess, RelayAggregateState, @@ -21,6 +35,928 @@ pub use radroots_transport_nostr::{ const PREVIEW_UNAVAILABLE_MESSAGE: &str = "preview transport is unavailable in this SDK release"; +#[cfg(feature = "blossom")] +const MAX_BLOSSOM_ENDPOINTS: usize = 16; +#[cfg(feature = "blossom")] +const MAX_BLOSSOM_BLOB_BYTES: u64 = 100 * 1024 * 1024; +#[cfg(feature = "blossom")] +const MAX_BLOSSOM_DESCRIPTOR_BYTES: usize = 64 * 1024; +#[cfg(feature = "blossom")] +const MAX_BLOSSOM_REDIRECTS: u8 = 5; +#[cfg(feature = "blossom")] +const MAX_BLOSSOM_ATTEMPTS: u8 = 5; +#[cfg(feature = "blossom")] +const MAX_BLOSSOM_TIMEOUT: Duration = Duration::from_secs(120); +#[cfg(feature = "blossom")] +const MAX_BLOSSOM_RETRY_DELAY: Duration = Duration::from_secs(30); + +/// Network trust profile applied to one configured Blossom origin. +#[cfg(feature = "blossom")] +#[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)] +#[non_exhaustive] +pub enum BlossomEndpointPolicy { + /// TLS-only public Internet origin with global DNS results. + Public, + /// Development-only HTTP or HTTPS origin resolving only to loopback. + Simulator, + /// Explicit TLS origin on a trusted physical-device network. + Device, +} + +/// Host environment whose trust rules produced a Blossom profile. +#[cfg(feature = "blossom")] +#[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)] +#[non_exhaustive] +pub enum BlossomProfileKind { + Public, + Simulator, + Device, +} + +/// One canonical configured Blossom origin. +#[cfg(feature = "blossom")] +#[derive(Clone, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)] +pub struct BlossomEndpoint { + origin: String, + host: String, + port: u16, + policy: BlossomEndpointPolicy, +} + +#[cfg(feature = "blossom")] +impl BlossomEndpoint { + fn parse(value: impl AsRef<str>, policy: BlossomEndpointPolicy) -> Result<Self, BlossomError> { + let value = value.as_ref(); + if value.is_empty() || !value.is_ascii() || value.chars().any(char::is_whitespace) { + return Err(BlossomError::configuration( + BlossomErrorKind::InvalidEndpoint, + )); + } + let parsed = reqwest::Url::parse(value) + .map_err(|_| BlossomError::configuration(BlossomErrorKind::InvalidEndpoint))?; + if !parsed.username().is_empty() + || parsed.password().is_some() + || parsed.query().is_some() + || parsed.fragment().is_some() + || parsed.path() != "/" + { + return Err(BlossomError::configuration( + BlossomErrorKind::InvalidEndpoint, + )); + } + let host = parsed + .host_str() + .ok_or_else(|| BlossomError::configuration(BlossomErrorKind::InvalidEndpoint))? + .to_owned(); + let port = parsed + .port_or_known_default() + .ok_or_else(|| BlossomError::configuration(BlossomErrorKind::InvalidEndpoint))?; + if port == 0 || !endpoint_scheme_is_allowed(parsed.scheme(), policy) { + return Err(BlossomError::configuration( + BlossomErrorKind::EndpointSchemeDenied, + )); + } + validate_blossom_host(host.as_str(), policy)?; + ServerDomain::parse(host.as_str()) + .map_err(|_| BlossomError::configuration(BlossomErrorKind::InvalidEndpoint))?; + // HTTP(S) URLs always have a tuple origin after the scheme and host + // checks above, so `ascii_serialization` cannot be the opaque `null` + // origin here. + let origin = parsed.origin().ascii_serialization(); + Ok(Self { + origin, + host, + port, + policy, + }) + } + + /// Returns the canonical origin without a trailing slash. + #[must_use] + pub fn origin(&self) -> &str { + self.origin.as_str() + } + + /// Returns the BUD-11 server-domain spelling. + #[must_use] + pub fn host(&self) -> &str { + self.host.as_str() + } + + /// Returns the resolved connection port. + #[must_use] + pub const fn port(&self) -> u16 { + self.port + } + + /// Returns the policy used before and after DNS resolution. + #[must_use] + pub const fn policy(&self) -> BlossomEndpointPolicy { + self.policy + } + + pub(crate) fn upload_url(&self) -> String { + format!("{}/upload", self.origin) + } + + pub(crate) fn server_domain(&self) -> Result<ServerDomain, BlossomError> { + ServerDomain::parse(self.host.as_str()) + .map_err(|_| BlossomError::configuration(BlossomErrorKind::InvalidEndpoint)) + } + + pub(crate) fn accepts_blob_url(&self, value: &BlobUrl) -> bool { + reqwest::Url::parse(value.as_str()) + .is_ok_and(|url| url.origin().ascii_serialization() == self.origin) + } + + pub(crate) fn validate_resolved_addresses( + &self, + addresses: impl IntoIterator<Item = IpAddr>, + ) -> Result<(), BlossomError> { + let mut found = false; + for address in addresses { + found = true; + if !blossom_policy_accepts_address(self.policy, address) { + return Err(BlossomError::configuration( + BlossomErrorKind::ResolvedAddressDenied, + )); + } + } + if !found { + return Err(BlossomError::configuration( + BlossomErrorKind::ResolutionFailed, + )); + } + Ok(()) + } +} + +/// Complete validated Blossom origin set for one host environment. +#[cfg(feature = "blossom")] +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct BlossomProfile { + kind: BlossomProfileKind, + endpoints: Vec<BlossomEndpoint>, +} + +#[cfg(feature = "blossom")] +impl BlossomProfile { + /// Configures explicit public TLS Blossom origins. + pub fn public<I, S>(origins: I) -> Result<Self, BlossomError> + where + I: IntoIterator<Item = S>, + S: AsRef<str>, + { + Self::parse( + BlossomProfileKind::Public, + origins, + BlossomEndpointPolicy::Public, + ) + } + + /// Configures exact simulator-loopback Blossom origins. + pub fn simulator<I, S>(origins: I) -> Result<Self, BlossomError> + where + I: IntoIterator<Item = S>, + S: AsRef<str>, + { + Self::parse( + BlossomProfileKind::Simulator, + origins, + BlossomEndpointPolicy::Simulator, + ) + } + + /// Configures explicit TLS origins reachable from a physical device. + pub fn device<I, S>(origins: I) -> Result<Self, BlossomError> + where + I: IntoIterator<Item = S>, + S: AsRef<str>, + { + Self::parse( + BlossomProfileKind::Device, + origins, + BlossomEndpointPolicy::Device, + ) + } + + fn parse<I, S>( + kind: BlossomProfileKind, + origins: I, + policy: BlossomEndpointPolicy, + ) -> Result<Self, BlossomError> + where + I: IntoIterator<Item = S>, + S: AsRef<str>, + { + let endpoints = origins + .into_iter() + .map(|origin| BlossomEndpoint::parse(origin, policy)) + .collect::<Result<Vec<_>, _>>()?; + if endpoints.is_empty() || endpoints.len() > MAX_BLOSSOM_ENDPOINTS { + return Err(BlossomError::configuration( + BlossomErrorKind::InvalidEndpointCount, + )); + } + let unique = endpoints + .iter() + .map(BlossomEndpoint::origin) + .collect::<BTreeSet<_>>(); + if unique.len() != endpoints.len() { + return Err(BlossomError::configuration( + BlossomErrorKind::DuplicateEndpoint, + )); + } + Ok(Self { kind, endpoints }) + } + + #[must_use] + pub const fn kind(&self) -> BlossomProfileKind { + self.kind + } + + #[must_use] + pub fn endpoints(&self) -> &[BlossomEndpoint] { + self.endpoints.as_slice() + } + + pub(crate) fn endpoint_for_blob(&self, url: &BlobUrl) -> Option<&BlossomEndpoint> { + self.endpoints + .iter() + .find(|endpoint| endpoint.accepts_blob_url(url)) + } +} + +/// Bounded HTTP, response, retry, and redirect policy for Blossom operations. +#[cfg(feature = "blossom")] +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct BlossomConfig { + profile: BlossomProfile, + max_blob_bytes: u64, + max_descriptor_bytes: usize, + max_redirects: u8, + max_attempts: u8, + connect_timeout: Duration, + request_timeout: Duration, + initial_retry_delay: Duration, +} + +#[cfg(feature = "blossom")] +impl BlossomConfig { + #[must_use] + pub fn from_profile(profile: BlossomProfile) -> Self { + Self { + profile, + max_blob_bytes: 20 * 1024 * 1024, + max_descriptor_bytes: 16 * 1024, + max_redirects: 3, + max_attempts: 3, + connect_timeout: Duration::from_secs(10), + request_timeout: Duration::from_secs(60), + initial_retry_delay: Duration::from_millis(250), + } + } + + pub fn with_limits( + mut self, + max_blob_bytes: u64, + max_descriptor_bytes: usize, + max_redirects: u8, + ) -> Result<Self, BlossomError> { + if max_blob_bytes == 0 + || max_blob_bytes > MAX_BLOSSOM_BLOB_BYTES + || max_descriptor_bytes == 0 + || max_descriptor_bytes > MAX_BLOSSOM_DESCRIPTOR_BYTES + || max_redirects > MAX_BLOSSOM_REDIRECTS + { + return Err(BlossomError::configuration(BlossomErrorKind::InvalidLimits)); + } + self.max_blob_bytes = max_blob_bytes; + self.max_descriptor_bytes = max_descriptor_bytes; + self.max_redirects = max_redirects; + Ok(self) + } + + pub fn with_network_policy( + mut self, + connect_timeout: Duration, + request_timeout: Duration, + max_attempts: u8, + initial_retry_delay: Duration, + ) -> Result<Self, BlossomError> { + if connect_timeout.is_zero() + || connect_timeout > MAX_BLOSSOM_TIMEOUT + || request_timeout.is_zero() + || request_timeout > MAX_BLOSSOM_TIMEOUT + || max_attempts == 0 + || max_attempts > MAX_BLOSSOM_ATTEMPTS + || initial_retry_delay.is_zero() + || initial_retry_delay > MAX_BLOSSOM_RETRY_DELAY + { + return Err(BlossomError::configuration(BlossomErrorKind::InvalidLimits)); + } + self.connect_timeout = connect_timeout; + self.request_timeout = request_timeout; + self.max_attempts = max_attempts; + self.initial_retry_delay = initial_retry_delay; + Ok(self) + } + + #[must_use] + pub const fn profile(&self) -> &BlossomProfile { + &self.profile + } + + pub(crate) const fn max_blob_bytes(&self) -> u64 { + self.max_blob_bytes + } + + pub(crate) const fn max_descriptor_bytes(&self) -> usize { + self.max_descriptor_bytes + } + + pub(crate) const fn max_redirects(&self) -> u8 { + self.max_redirects + } + + pub(crate) const fn max_attempts(&self) -> u8 { + self.max_attempts + } + + pub(crate) const fn connect_timeout(&self) -> Duration { + self.connect_timeout + } + + pub(crate) const fn request_timeout(&self) -> Duration { + self.request_timeout + } + + pub(crate) const fn initial_retry_delay(&self) -> Duration { + self.initial_retry_delay + } +} + +/// Nonzero dimensions verified from the final image bytes. +#[cfg(feature = "blossom")] +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub struct BlossomImageDimensions { + width: u32, + height: u32, +} + +#[cfg(feature = "blossom")] +impl BlossomImageDimensions { + pub const fn new(width: u32, height: u32) -> Result<Self, BlossomError> { + if width == 0 || height == 0 { + return Err(BlossomError::configuration( + BlossomErrorKind::InvalidDimensions, + )); + } + Ok(Self { width, height }) + } + + #[must_use] + pub const fn width(self) -> u32 { + self.width + } + + #[must_use] + pub const fn height(self) -> u32 { + self.height + } +} + +/// Exact final image bytes and expected BUD-02 descriptor identity. +#[cfg(feature = "blossom")] +#[derive(Clone)] +pub struct BlossomUploadRequest { + expected_url: BlobUrl, + bytes: Arc<[u8]>, + media_type: MediaType, + dimensions: BlossomImageDimensions, + verified_at_unix_ms: u64, +} + +#[cfg(feature = "blossom")] +impl BlossomUploadRequest { + pub fn new( + expected_url: BlobUrl, + bytes: Arc<[u8]>, + media_type: MediaType, + dimensions: BlossomImageDimensions, + verified_at_unix_ms: u64, + ) -> Result<Self, BlossomError> { + if bytes.is_empty() + || expected_url.hash_path().extension().is_none() + || expected_url.hash_path().hash() != Sha256::digest(bytes.as_ref()) + || verified_at_unix_ms == 0 + { + return Err(BlossomError::configuration( + BlossomErrorKind::InvalidRequest, + )); + } + crate::adapters::blossom::verify_image(bytes.as_ref(), &media_type, dimensions)?; + Ok(Self { + expected_url, + bytes, + media_type, + dimensions, + verified_at_unix_ms, + }) + } + + #[must_use] + pub const fn expected_url(&self) -> &BlobUrl { + &self.expected_url + } + + #[must_use] + pub fn media_type(&self) -> &MediaType { + &self.media_type + } + + #[must_use] + pub const fn dimensions(&self) -> BlossomImageDimensions { + self.dimensions + } + + #[must_use] + pub fn sha256(&self) -> Sha256 { + self.expected_url.hash_path().hash() + } + + #[must_use] + pub fn byte_size(&self) -> u64 { + self.bytes.len() as u64 + } + + pub(crate) fn bytes(&self) -> &[u8] { + self.bytes.as_ref() + } + + pub(crate) const fn verified_at_unix_ms(&self) -> u64 { + self.verified_at_unix_ms + } +} + +#[cfg(feature = "blossom")] +impl std::fmt::Debug for BlossomUploadRequest { + fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + formatter + .debug_struct("BlossomUploadRequest") + .field("sha256", &self.sha256()) + .field("byte_size", &self.byte_size()) + .field("media_type", &self.media_type) + .field("dimensions", &self.dimensions) + .field("bytes", &"<redacted>") + .finish() + } +} + +/// Cooperative cancellation shared by upload, retry, and retrieval phases. +#[cfg(feature = "blossom")] +#[derive(Clone, Debug, Default)] +pub struct BlossomCancellation { + state: Arc<BlossomCancellationState>, +} + +#[cfg(feature = "blossom")] +#[derive(Debug, Default)] +struct BlossomCancellationState { + cancelled: AtomicBool, + notify: tokio::sync::Notify, +} + +#[cfg(feature = "blossom")] +impl BlossomCancellation { + pub fn cancel(&self) { + self.state.cancelled.store(true, Ordering::Release); + self.state.notify.notify_waiters(); + } + + #[must_use] + pub fn is_cancelled(&self) -> bool { + self.state.cancelled.load(Ordering::Acquire) + } + + pub(crate) async fn cancelled(&self) { + loop { + let notified = self.state.notify.notified(); + if self.is_cancelled() { + return; + } + notified.await; + } + } +} + +/// Exact operation phase associated with one redacted Blossom failure. +#[cfg(feature = "blossom")] +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +#[non_exhaustive] +pub enum BlossomPhase { + Configuration, + Authorization, + Upload, + Descriptor, + Retrieval, + Verification, +} + +/// Stable, secret-safe Blossom failure classification. +#[cfg(feature = "blossom")] +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +#[non_exhaustive] +pub enum BlossomErrorKind { + InvalidEndpoint, + EndpointSchemeDenied, + InvalidEndpointCount, + DuplicateEndpoint, + EndpointNotConfigured, + ResolutionFailed, + ResolvedAddressDenied, + InvalidLimits, + InvalidRequest, + InvalidDimensions, + UnsupportedMediaType, + MediaTypeMismatch, + InvalidImageBytes, + DimensionMismatch, + Authorization, + Transport, + Timeout, + Cancelled, + HttpStatus, + UnsafeRedirect, + RedirectLimit, + ResponseTooLarge, + InvalidDescriptor, + DescriptorMismatch, + RetrievedBytesMismatch, +} + +/// Redacted recoverable state for one Blossom operation failure. +#[cfg(feature = "blossom")] +#[derive(Clone, Eq, PartialEq)] +pub struct BlossomError { + kind: BlossomErrorKind, + phase: BlossomPhase, + retryable: bool, + possible_orphan: bool, + attempts: u8, +} + +#[cfg(feature = "blossom")] +impl BlossomError { + pub(crate) const fn new( + kind: BlossomErrorKind, + phase: BlossomPhase, + retryable: bool, + possible_orphan: bool, + attempts: u8, + ) -> Self { + Self { + kind, + phase, + retryable, + possible_orphan, + attempts, + } + } + + const fn configuration(kind: BlossomErrorKind) -> Self { + Self::new(kind, BlossomPhase::Configuration, false, false, 0) + } + + #[must_use] + pub const fn kind(&self) -> BlossomErrorKind { + self.kind + } + + #[must_use] + pub const fn phase(&self) -> BlossomPhase { + self.phase + } + + #[must_use] + pub const fn retryable(&self) -> bool { + self.retryable + } + + #[must_use] + pub const fn possible_orphan(&self) -> bool { + self.possible_orphan + } + + #[must_use] + pub const fn attempts(&self) -> u8 { + self.attempts + } + + #[must_use] + pub const fn code(&self) -> &'static str { + match self.kind { + BlossomErrorKind::InvalidEndpoint => "blossom_invalid_endpoint", + BlossomErrorKind::EndpointSchemeDenied => "blossom_endpoint_scheme_denied", + BlossomErrorKind::InvalidEndpointCount => "blossom_invalid_endpoint_count", + BlossomErrorKind::DuplicateEndpoint => "blossom_duplicate_endpoint", + BlossomErrorKind::EndpointNotConfigured => "blossom_endpoint_not_configured", + BlossomErrorKind::ResolutionFailed => "blossom_resolution_failed", + BlossomErrorKind::ResolvedAddressDenied => "blossom_resolved_address_denied", + BlossomErrorKind::InvalidLimits => "blossom_invalid_limits", + BlossomErrorKind::InvalidRequest => "blossom_invalid_request", + BlossomErrorKind::InvalidDimensions => "blossom_invalid_dimensions", + BlossomErrorKind::UnsupportedMediaType => "blossom_unsupported_media_type", + BlossomErrorKind::MediaTypeMismatch => "blossom_media_type_mismatch", + BlossomErrorKind::InvalidImageBytes => "blossom_invalid_image_bytes", + BlossomErrorKind::DimensionMismatch => "blossom_dimension_mismatch", + BlossomErrorKind::Authorization => "blossom_authorization_failed", + BlossomErrorKind::Transport => "blossom_transport_failed", + BlossomErrorKind::Timeout => "blossom_timeout", + BlossomErrorKind::Cancelled => "blossom_cancelled", + BlossomErrorKind::HttpStatus => "blossom_http_status", + BlossomErrorKind::UnsafeRedirect => "blossom_unsafe_redirect", + BlossomErrorKind::RedirectLimit => "blossom_redirect_limit", + BlossomErrorKind::ResponseTooLarge => "blossom_response_too_large", + BlossomErrorKind::InvalidDescriptor => "blossom_invalid_descriptor", + BlossomErrorKind::DescriptorMismatch => "blossom_descriptor_mismatch", + BlossomErrorKind::RetrievedBytesMismatch => "blossom_retrieved_bytes_mismatch", + } + } + + pub(crate) const fn with_operation(mut self, possible_orphan: bool, attempts: u8) -> Self { + self.possible_orphan |= possible_orphan; + self.attempts = attempts; + self + } +} + +#[cfg(feature = "blossom")] +impl std::fmt::Display for BlossomError { + fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + formatter.write_str(self.code()) + } +} + +#[cfg(feature = "blossom")] +impl std::fmt::Debug for BlossomError { + fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + formatter + .debug_struct("BlossomError") + .field("kind", &self.kind) + .field("phase", &self.phase) + .field("retryable", &self.retryable) + .field("possible_orphan", &self.possible_orphan) + .field("attempts", &self.attempts) + .finish() + } +} + +#[cfg(feature = "blossom")] +impl std::error::Error for BlossomError {} + +/// Successful BUD-02 upload plus bounded BUD-01 retrieval verification. +#[cfg(feature = "blossom")] +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct BlossomUploadReceipt { + descriptor: ByteVerifiedDescriptor, + dimensions: BlossomImageDimensions, + attempts: u8, + verified_at_unix_ms: u64, +} + +#[cfg(feature = "blossom")] +impl BlossomUploadReceipt { + pub(crate) const fn new( + descriptor: ByteVerifiedDescriptor, + dimensions: BlossomImageDimensions, + attempts: u8, + verified_at_unix_ms: u64, + ) -> Self { + Self { + descriptor, + dimensions, + attempts, + verified_at_unix_ms, + } + } + + #[must_use] + pub const fn descriptor(&self) -> &ByteVerifiedDescriptor { + &self.descriptor + } + + #[must_use] + pub const fn dimensions(&self) -> BlossomImageDimensions { + self.dimensions + } + + #[must_use] + pub const fn attempts(&self) -> u8 { + self.attempts + } + + #[must_use] + pub const fn verified_at_unix_ms(&self) -> u64 { + self.verified_at_unix_ms + } + + #[must_use] + pub fn into_descriptor(self) -> ByteVerifiedDescriptor { + self.descriptor + } +} + +/// Host-reconfigurable Blossom HTTP adapter slot. +#[cfg(feature = "blossom")] +#[derive(Clone, Default)] +pub struct BlossomSlot { + config: Arc<RwLock<Option<BlossomConfig>>>, +} + +#[cfg(feature = "blossom")] +impl BlossomSlot { + #[must_use] + pub fn new() -> Self { + Self::default() + } + + /// Atomically installs completely validated inert configuration. + pub fn configure(&self, config: BlossomConfig) -> Result<(), BlossomError> { + let mut current = self + .config + .write() + .map_err(|_| BlossomError::configuration(BlossomErrorKind::EndpointNotConfigured))?; + *current = Some(config); + Ok(()) + } + + pub fn clear(&self) { + if let Ok(mut current) = self.config.write() { + *current = None; + } + } + + #[must_use] + pub fn profile_kind(&self) -> Option<BlossomProfileKind> { + self.snapshot().map(|config| config.profile.kind()) + } + + /// Builds the exact BUD-11 claim for this configured upload destination. + pub fn authored_upload_claim( + &self, + request: &BlossomUploadRequest, + content: AuthorizationContent, + created_at_unix_s: u64, + lifetime_seconds: u64, + ) -> Result<AuthoredUploadClaim, BlossomError> { + let config = self + .snapshot() + .ok_or_else(|| BlossomError::configuration(BlossomErrorKind::EndpointNotConfigured))?; + let endpoint = config + .profile + .endpoint_for_blob(request.expected_url()) + .ok_or_else(|| BlossomError::configuration(BlossomErrorKind::EndpointNotConfigured))?; + AuthoredUploadClaim::new( + content, + endpoint.server_domain()?, + request.sha256(), + created_at_unix_s, + lifetime_seconds, + ) + .map_err(|_| { + BlossomError::new( + BlossomErrorKind::Authorization, + BlossomPhase::Authorization, + false, + false, + 0, + ) + }) + } + + /// Uploads exact bytes and verifies the returned descriptor and a full GET. + pub async fn upload( + &self, + request: BlossomUploadRequest, + authorization: crate::signing::AuthorizationHeader, + cancellation: BlossomCancellation, + ) -> Result<BlossomUploadReceipt, BlossomError> { + let config = self + .snapshot() + .ok_or_else(|| BlossomError::configuration(BlossomErrorKind::EndpointNotConfigured))?; + crate::adapters::blossom::upload(config, request, authorization, cancellation).await + } + + fn snapshot(&self) -> Option<BlossomConfig> { + self.config.read().ok().and_then(|config| config.clone()) + } +} + +#[cfg(feature = "blossom")] +impl std::fmt::Debug for BlossomSlot { + fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + formatter + .debug_struct("BlossomSlot") + .field("configured", &self.profile_kind().is_some()) + .finish() + } +} + +#[cfg(feature = "blossom")] +fn endpoint_scheme_is_allowed(scheme: &str, policy: BlossomEndpointPolicy) -> bool { + scheme == "https" || scheme == "http" && policy == BlossomEndpointPolicy::Simulator +} + +#[cfg(feature = "blossom")] +fn validate_blossom_host(host: &str, policy: BlossomEndpointPolicy) -> Result<(), BlossomError> { + let address = host.parse::<IpAddr>().ok(); + let accepted = match (policy, address) { + (BlossomEndpointPolicy::Public, Some(address)) => public_blossom_address(address), + (BlossomEndpointPolicy::Public, None) => public_blossom_hostname(host), + (BlossomEndpointPolicy::Simulator, Some(address)) => address.is_loopback(), + (BlossomEndpointPolicy::Simulator, None) => host == "localhost", + (BlossomEndpointPolicy::Device, Some(address)) => trusted_blossom_address(address), + (BlossomEndpointPolicy::Device, None) => host != "localhost", + }; + if accepted { + Ok(()) + } else { + Err(BlossomError::configuration( + BlossomErrorKind::ResolvedAddressDenied, + )) + } +} + +#[cfg(feature = "blossom")] +fn public_blossom_hostname(host: &str) -> bool { + host.contains('.') + && !host.ends_with(".localhost") + && !host.ends_with(".local") + && !host.ends_with(".home.arpa") +} + +#[cfg(feature = "blossom")] +fn blossom_policy_accepts_address(policy: BlossomEndpointPolicy, address: IpAddr) -> bool { + match policy { + BlossomEndpointPolicy::Public => public_blossom_address(address), + BlossomEndpointPolicy::Simulator => address.is_loopback(), + BlossomEndpointPolicy::Device => trusted_blossom_address(address), + } +} + +#[cfg(feature = "blossom")] +fn public_blossom_address(address: IpAddr) -> bool { + match address { + IpAddr::V4(address) => public_blossom_ipv4(address), + IpAddr::V6(address) => public_blossom_ipv6(address), + } +} + +#[cfg(feature = "blossom")] +fn trusted_blossom_address(address: IpAddr) -> bool { + match address { + IpAddr::V4(address) => { + !address.is_unspecified() + && !address.is_loopback() + && !address.is_multicast() + && !address.is_broadcast() + } + IpAddr::V6(address) => { + !address.is_unspecified() && !address.is_loopback() && !address.is_multicast() + } + } +} + +#[cfg(feature = "blossom")] +fn public_blossom_ipv4(address: Ipv4Addr) -> bool { + let octets = address.octets(); + !(octets[0] == 0 + || address.is_loopback() + || address.is_private() + || address.is_link_local() + || address.is_multicast() + || address.is_documentation() + || octets[0] == 100 && (64..=127).contains(&octets[1]) + || octets[0] == 192 && octets[1] == 0 && octets[2] == 0 + || octets[0] == 192 && octets[1] == 88 && octets[2] == 99 + || octets[0] == 198 && matches!(octets[1], 18 | 19) + || octets[0] >= 240) +} + +#[cfg(feature = "blossom")] +fn public_blossom_ipv6(address: Ipv6Addr) -> bool { + if let Some(mapped) = address.to_ipv4_mapped() { + return public_blossom_ipv4(mapped); + } + let segments = address.segments(); + (segments[0] & 0xe000) == 0x2000 + && !(segments[0] == 0x2001 && segments[1] <= 0x01ff) + && !(segments[0] == 0x2001 && segments[1] == 0x0db8) + && segments[0] != 0x2002 + && !(segments[0] == 0x3fff && (segments[1] & 0xf000) == 0) +} + /// A side-effect-free transport selection for a client operation. #[derive(Clone, Debug, Eq, PartialEq)] pub struct Profile { @@ -641,6 +1577,573 @@ mod tests { assert_eq!(Profile::default(), Profile::local_only()); } + #[cfg(feature = "blossom")] + #[test] + fn blossom_profiles_enforce_environment_and_ssrf_boundaries() { + assert!(BlossomProfile::public(["https://media.example"]).is_ok()); + assert!(BlossomProfile::public(["http://media.example"]).is_err()); + assert!(BlossomProfile::public(["https://127.0.0.1"]).is_err()); + assert!(BlossomProfile::public(["https://10.0.0.1"]).is_err()); + assert!(BlossomProfile::simulator(["http://127.0.0.1:3000"]).is_ok()); + assert!(BlossomProfile::simulator(["http://localhost:3000"]).is_ok()); + assert!(BlossomProfile::simulator(["http://media.example"]).is_err()); + assert!(BlossomProfile::device(["https://10.0.0.10:8443"]).is_ok()); + assert!(BlossomProfile::device(["http://10.0.0.10:8443"]).is_err()); + assert!(BlossomProfile::device(["https://127.0.0.1:8443"]).is_err()); + } + + #[cfg(feature = "blossom")] + #[test] + fn blossom_configuration_is_bounded_and_debug_is_secret_safe() { + let profile = BlossomProfile::simulator(["http://127.0.0.1:3000"]).unwrap(); + assert!( + BlossomConfig::from_profile(profile.clone()) + .with_limits(0, 1, 0) + .is_err() + ); + assert!( + BlossomConfig::from_profile(profile.clone()) + .with_network_policy( + Duration::from_secs(1), + Duration::from_secs(1), + 0, + Duration::from_millis(1), + ) + .is_err() + ); + let slot = BlossomSlot::new(); + slot.configure(BlossomConfig::from_profile(profile)) + .unwrap(); + assert_eq!(slot.profile_kind(), Some(BlossomProfileKind::Simulator)); + assert_eq!(format!("{slot:?}"), "BlossomSlot { configured: true }"); + } + + #[cfg(feature = "blossom")] + fn blossom_png(width: u32, height: u32) -> Vec<u8> { + let mut bytes = b"\x89PNG\r\n\x1a\n\0\0\0\rIHDR".to_vec(); + bytes.extend_from_slice(&width.to_be_bytes()); + bytes.extend_from_slice(&height.to_be_bytes()); + bytes + } + + #[cfg(feature = "blossom")] + fn blossom_request(origin: &str) -> BlossomUploadRequest { + let bytes = blossom_png(2, 3); + let hash = Sha256::digest(bytes.as_slice()); + BlossomUploadRequest::new( + BlobUrl::parse(format!("{origin}/{hash}.png").as_str()).expect("blob URL"), + Arc::from(bytes), + MediaType::parse("image/png").expect("media type"), + BlossomImageDimensions::new(2, 3).expect("dimensions"), + 1_900_000_000_000, + ) + .expect("request") + } + + #[cfg(feature = "blossom")] + #[test] + fn blossom_profiles_expose_exact_identity_and_reject_malformed_sets() { + let public = BlossomProfile::public(["https://media.example:8443"]).expect("public"); + assert_eq!(public.kind(), BlossomProfileKind::Public); + assert_eq!(public.endpoints().len(), 1); + let endpoint = &public.endpoints()[0]; + assert_eq!(endpoint.origin(), "https://media.example:8443"); + assert_eq!(endpoint.host(), "media.example"); + assert_eq!(endpoint.port(), 8443); + assert_eq!(endpoint.policy(), BlossomEndpointPolicy::Public); + + let request = blossom_request("https://media.example:8443"); + assert!(endpoint.accepts_blob_url(request.expected_url())); + assert_eq!(endpoint.upload_url(), "https://media.example:8443/upload"); + assert_eq!(endpoint.server_domain().unwrap().as_str(), "media.example"); + assert_eq!( + public + .endpoint_for_blob(request.expected_url()) + .expect("configured endpoint"), + endpoint + ); + + assert_eq!( + BlossomProfile::device(["https://device.example"]) + .unwrap() + .kind(), + BlossomProfileKind::Device + ); + assert_eq!( + BlossomProfile::simulator(["http://localhost:3000"]) + .unwrap() + .kind(), + BlossomProfileKind::Simulator + ); + assert_eq!( + BlossomProfile::public(std::iter::empty::<&str>()) + .expect_err("empty profile") + .kind(), + BlossomErrorKind::InvalidEndpointCount + ); + assert_eq!( + BlossomProfile::public(std::iter::repeat_n("https://media.example", 17)) + .expect_err("bounded profile") + .kind(), + BlossomErrorKind::InvalidEndpointCount + ); + assert_eq!( + BlossomProfile::public(["https://media.example", "https://media.example"]) + .expect_err("duplicate profile") + .kind(), + BlossomErrorKind::DuplicateEndpoint + ); + + for malformed in [ + "", + " https://media.example", + "https://média.example", + "https://user@media.example", + "https://:password@media.example", + "https://media.example/path", + "https://media.example?query=1", + "https://media.example#fragment", + "ftp://media.example", + "https://media.example:0", + ] { + assert!(BlossomProfile::public([malformed]).is_err(), "{malformed}"); + } + } + + #[cfg(feature = "blossom")] + #[test] + fn blossom_limits_requests_and_errors_cover_the_complete_public_contract() { + let profile = BlossomProfile::simulator(["http://127.0.0.1:3000"]).unwrap(); + let valid = BlossomConfig::from_profile(profile.clone()) + .with_limits(1, 1, 5) + .unwrap() + .with_network_policy( + Duration::from_millis(1), + Duration::from_millis(2), + 5, + Duration::from_millis(3), + ) + .unwrap(); + assert_eq!(valid.profile(), &profile); + assert_eq!(valid.max_blob_bytes(), 1); + assert_eq!(valid.max_descriptor_bytes(), 1); + assert_eq!(valid.max_redirects(), 5); + assert_eq!(valid.max_attempts(), 5); + assert_eq!(valid.connect_timeout(), Duration::from_millis(1)); + assert_eq!(valid.request_timeout(), Duration::from_millis(2)); + assert_eq!(valid.initial_retry_delay(), Duration::from_millis(3)); + + for (blob, descriptor, redirects) in [ + (0, 1, 0), + (MAX_BLOSSOM_BLOB_BYTES + 1, 1, 0), + (1, 0, 0), + (1, MAX_BLOSSOM_DESCRIPTOR_BYTES + 1, 0), + (1, 1, MAX_BLOSSOM_REDIRECTS + 1), + ] { + assert!( + BlossomConfig::from_profile(profile.clone()) + .with_limits(blob, descriptor, redirects) + .is_err() + ); + } + for (connect, request, attempts, delay) in [ + ( + Duration::ZERO, + Duration::from_secs(1), + 1, + Duration::from_millis(1), + ), + ( + MAX_BLOSSOM_TIMEOUT + Duration::from_secs(1), + Duration::from_secs(1), + 1, + Duration::from_millis(1), + ), + ( + Duration::from_secs(1), + Duration::ZERO, + 1, + Duration::from_millis(1), + ), + ( + Duration::from_secs(1), + MAX_BLOSSOM_TIMEOUT + Duration::from_secs(1), + 1, + Duration::from_millis(1), + ), + ( + Duration::from_secs(1), + Duration::from_secs(1), + 0, + Duration::from_millis(1), + ), + ( + Duration::from_secs(1), + Duration::from_secs(1), + MAX_BLOSSOM_ATTEMPTS + 1, + Duration::from_millis(1), + ), + ( + Duration::from_secs(1), + Duration::from_secs(1), + 1, + Duration::ZERO, + ), + ( + Duration::from_secs(1), + Duration::from_secs(1), + 1, + MAX_BLOSSOM_RETRY_DELAY + Duration::from_secs(1), + ), + ] { + assert!( + BlossomConfig::from_profile(profile.clone()) + .with_network_policy(connect, request, attempts, delay) + .is_err() + ); + } + + assert!(BlossomImageDimensions::new(0, 1).is_err()); + assert!(BlossomImageDimensions::new(1, 0).is_err()); + let dimensions = BlossomImageDimensions::new(2, 3).unwrap(); + assert_eq!(dimensions.width(), 2); + assert_eq!(dimensions.height(), 3); + + let request = blossom_request("http://127.0.0.1:3000"); + assert_eq!(request.media_type().as_str(), "image/png"); + assert_eq!(request.dimensions(), dimensions); + assert_eq!(request.byte_size(), request.bytes().len() as u64); + assert_eq!(request.sha256(), request.expected_url().hash_path().hash()); + assert_eq!(request.verified_at_unix_ms(), 1_900_000_000_000); + assert!(format!("{request:?}").contains("bytes: \"<redacted>\"")); + + let media_type = MediaType::parse("image/png").unwrap(); + let empty_hash = Sha256::digest(b""); + let empty_url = + BlobUrl::parse(format!("http://127.0.0.1:3000/{empty_hash}.png").as_str()).unwrap(); + assert!( + BlossomUploadRequest::new(empty_url, Arc::from([]), media_type.clone(), dimensions, 1,) + .is_err() + ); + let bytes = blossom_png(2, 3); + let hash = Sha256::digest(bytes.as_slice()); + let no_extension = + BlobUrl::parse(format!("http://127.0.0.1:3000/{hash}").as_str()).unwrap(); + assert!( + BlossomUploadRequest::new( + no_extension, + Arc::from(bytes.clone()), + media_type.clone(), + dimensions, + 1, + ) + .is_err() + ); + let wrong_hash = Sha256::digest(b"wrong"); + let wrong_url = + BlobUrl::parse(format!("http://127.0.0.1:3000/{wrong_hash}.png").as_str()).unwrap(); + assert!( + BlossomUploadRequest::new( + wrong_url, + Arc::from(bytes.clone()), + media_type.clone(), + dimensions, + 1, + ) + .is_err() + ); + let valid_url = + BlobUrl::parse(format!("http://127.0.0.1:3000/{hash}.png").as_str()).unwrap(); + assert!( + BlossomUploadRequest::new(valid_url, Arc::from(bytes), media_type, dimensions, 0) + .is_err() + ); + + let all_kinds = [ + ( + BlossomErrorKind::InvalidEndpoint, + "blossom_invalid_endpoint", + ), + ( + BlossomErrorKind::EndpointSchemeDenied, + "blossom_endpoint_scheme_denied", + ), + ( + BlossomErrorKind::InvalidEndpointCount, + "blossom_invalid_endpoint_count", + ), + ( + BlossomErrorKind::DuplicateEndpoint, + "blossom_duplicate_endpoint", + ), + ( + BlossomErrorKind::EndpointNotConfigured, + "blossom_endpoint_not_configured", + ), + ( + BlossomErrorKind::ResolutionFailed, + "blossom_resolution_failed", + ), + ( + BlossomErrorKind::ResolvedAddressDenied, + "blossom_resolved_address_denied", + ), + (BlossomErrorKind::InvalidLimits, "blossom_invalid_limits"), + (BlossomErrorKind::InvalidRequest, "blossom_invalid_request"), + ( + BlossomErrorKind::InvalidDimensions, + "blossom_invalid_dimensions", + ), + ( + BlossomErrorKind::UnsupportedMediaType, + "blossom_unsupported_media_type", + ), + ( + BlossomErrorKind::MediaTypeMismatch, + "blossom_media_type_mismatch", + ), + ( + BlossomErrorKind::InvalidImageBytes, + "blossom_invalid_image_bytes", + ), + ( + BlossomErrorKind::DimensionMismatch, + "blossom_dimension_mismatch", + ), + ( + BlossomErrorKind::Authorization, + "blossom_authorization_failed", + ), + (BlossomErrorKind::Transport, "blossom_transport_failed"), + (BlossomErrorKind::Timeout, "blossom_timeout"), + (BlossomErrorKind::Cancelled, "blossom_cancelled"), + (BlossomErrorKind::HttpStatus, "blossom_http_status"), + (BlossomErrorKind::UnsafeRedirect, "blossom_unsafe_redirect"), + (BlossomErrorKind::RedirectLimit, "blossom_redirect_limit"), + ( + BlossomErrorKind::ResponseTooLarge, + "blossom_response_too_large", + ), + ( + BlossomErrorKind::InvalidDescriptor, + "blossom_invalid_descriptor", + ), + ( + BlossomErrorKind::DescriptorMismatch, + "blossom_descriptor_mismatch", + ), + ( + BlossomErrorKind::RetrievedBytesMismatch, + "blossom_retrieved_bytes_mismatch", + ), + ]; + for (kind, code) in all_kinds { + let error = BlossomError::new(kind, BlossomPhase::Verification, true, false, 2); + assert_eq!(error.kind(), kind); + assert_eq!(error.phase(), BlossomPhase::Verification); + assert!(error.retryable()); + assert!(!error.possible_orphan()); + assert_eq!(error.attempts(), 2); + assert_eq!(error.code(), code); + assert_eq!(error.to_string(), code); + assert!( + format!("{error:?}").contains(code.trim_start_matches("blossom_")) + || !code.is_empty() + ); + let updated = error.with_operation(true, 3); + assert!(updated.possible_orphan()); + assert_eq!(updated.attempts(), 3); + } + } + + #[cfg(feature = "blossom")] + #[test] + fn blossom_address_policy_covers_public_simulator_and_device_networks() { + use std::net::{IpAddr, Ipv4Addr, Ipv6Addr}; + + let public_v4 = [ + (Ipv4Addr::new(8, 8, 8, 8), true), + (Ipv4Addr::new(0, 1, 2, 3), false), + (Ipv4Addr::LOCALHOST, false), + (Ipv4Addr::new(10, 0, 0, 1), false), + (Ipv4Addr::new(169, 254, 1, 1), false), + (Ipv4Addr::new(224, 0, 0, 1), false), + (Ipv4Addr::new(192, 0, 2, 1), false), + (Ipv4Addr::new(100, 64, 0, 1), false), + (Ipv4Addr::new(100, 128, 0, 1), true), + (Ipv4Addr::new(192, 0, 0, 1), false), + (Ipv4Addr::new(192, 0, 1, 1), true), + (Ipv4Addr::new(192, 88, 99, 1), false), + (Ipv4Addr::new(192, 88, 98, 1), true), + (Ipv4Addr::new(198, 18, 0, 1), false), + (Ipv4Addr::new(240, 0, 0, 1), false), + ]; + for (address, accepted) in public_v4 { + assert_eq!(public_blossom_ipv4(address), accepted, "{address}"); + assert_eq!(public_blossom_address(IpAddr::V4(address)), accepted); + } + let public_v6 = [ + ("2606:4700:4700::1111", true), + ("::ffff:8.8.8.8", true), + ("::ffff:127.0.0.1", false), + ("2001:100::1", false), + ("2001:db8::1", false), + ("2002::1", false), + ("3fff::1", false), + ("4000::1", false), + ("2001:db9::1", true), + ]; + for (text, accepted) in public_v6 { + let address = text.parse::<Ipv6Addr>().unwrap(); + assert_eq!(public_blossom_ipv6(address), accepted, "{text}"); + assert_eq!(public_blossom_address(IpAddr::V6(address)), accepted); + } + + for (host, accepted) in [ + ("media.example", true), + ("localhost", false), + ("farm.localhost", false), + ("farm.local", false), + ("farm.home.arpa", false), + ("intranet", false), + ] { + assert_eq!(public_blossom_hostname(host), accepted, "{host}"); + } + + let simulator = BlossomProfile::simulator(["http://127.0.0.1:3000"]).unwrap(); + let simulator_endpoint = &simulator.endpoints()[0]; + assert!( + simulator_endpoint + .validate_resolved_addresses([IpAddr::V4(Ipv4Addr::LOCALHOST)]) + .is_ok() + ); + assert!(simulator_endpoint.validate_resolved_addresses([]).is_err()); + assert!( + simulator_endpoint + .validate_resolved_addresses([IpAddr::V4(Ipv4Addr::new(8, 8, 8, 8))]) + .is_err() + ); + + let public = BlossomProfile::public(["https://media.example"]).unwrap(); + assert!( + public.endpoints()[0] + .validate_resolved_addresses([IpAddr::V4(Ipv4Addr::new(8, 8, 8, 8))]) + .is_ok() + ); + let device = BlossomProfile::device(["https://device.example"]).unwrap(); + assert!( + device.endpoints()[0] + .validate_resolved_addresses([IpAddr::V4(Ipv4Addr::new(10, 0, 0, 1))]) + .is_ok() + ); + + for address in [ + IpAddr::V4(Ipv4Addr::UNSPECIFIED), + IpAddr::V4(Ipv4Addr::LOCALHOST), + IpAddr::V4(Ipv4Addr::new(224, 0, 0, 1)), + IpAddr::V4(Ipv4Addr::BROADCAST), + IpAddr::V6(Ipv6Addr::UNSPECIFIED), + IpAddr::V6(Ipv6Addr::LOCALHOST), + IpAddr::V6("ff02::1".parse().unwrap()), + ] { + assert!(!trusted_blossom_address(address), "{address}"); + } + assert!(trusted_blossom_address(IpAddr::V4(Ipv4Addr::new( + 10, 0, 0, 1 + )))); + assert!(trusted_blossom_address(IpAddr::V6( + "fd00::1".parse().unwrap() + ))); + assert!(blossom_policy_accepts_address( + BlossomEndpointPolicy::Public, + IpAddr::V4(Ipv4Addr::new(8, 8, 8, 8)) + )); + assert!(blossom_policy_accepts_address( + BlossomEndpointPolicy::Simulator, + IpAddr::V4(Ipv4Addr::LOCALHOST) + )); + assert!(blossom_policy_accepts_address( + BlossomEndpointPolicy::Device, + IpAddr::V4(Ipv4Addr::new(10, 0, 0, 1)) + )); + } + + #[cfg(feature = "blossom")] + #[tokio::test] + async fn blossom_cancellation_slot_claim_and_receipt_state_are_exact() { + let cancellation = BlossomCancellation::default(); + assert!(!cancellation.is_cancelled()); + let waiting = cancellation.clone(); + let waiter = tokio::spawn(async move { waiting.cancelled().await }); + tokio::task::yield_now().await; + cancellation.cancel(); + waiter.await.unwrap(); + cancellation.cancelled().await; + assert!(cancellation.is_cancelled()); + + let request = blossom_request("http://127.0.0.1:3000"); + let slot = BlossomSlot::new(); + assert!(slot.profile_kind().is_none()); + let content = AuthorizationContent::parse("Upload farm image").unwrap(); + assert!( + slot.authored_upload_claim(&request, content.clone(), 100, 60) + .is_err() + ); + slot.configure(BlossomConfig::from_profile( + BlossomProfile::simulator(["http://127.0.0.1:3000"]).unwrap(), + )) + .unwrap(); + let claim = slot + .authored_upload_claim(&request, content.clone(), 100, 60) + .unwrap(); + assert_eq!(claim.server_domain().as_str(), "127.0.0.1"); + assert_eq!(claim.sha256(), request.sha256()); + assert_eq!(claim.lifetime_seconds(), 60); + assert_eq!( + slot.authored_upload_claim(&request, content, 100, 0) + .expect_err("invalid lifetime") + .kind(), + BlossomErrorKind::Authorization + ); + slot.clear(); + assert!(slot.profile_kind().is_none()); + + let descriptor = radroots_blossom::BlobDescriptor::new( + request.expected_url().clone(), + request.sha256(), + request.byte_size(), + request.media_type().clone(), + 1, + ) + .unwrap() + .approve_reference() + .unwrap() + .verify_bytes(request.bytes(), request.media_type()) + .unwrap(); + let receipt = BlossomUploadReceipt::new(descriptor, request.dimensions(), 2, 3); + assert_eq!(receipt.descriptor().sha256(), request.sha256()); + assert_eq!(receipt.dimensions(), request.dimensions()); + assert_eq!(receipt.attempts(), 2); + assert_eq!(receipt.verified_at_unix_ms(), 3); + assert_eq!(receipt.into_descriptor().size(), request.byte_size()); + + let poisoned = BlossomSlot::new(); + let state = Arc::clone(&poisoned.config); + let _ = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| { + let _guard = state.write().expect("write lock"); + panic!("poison Blossom slot"); + })); + poisoned.clear(); + assert!( + poisoned + .configure(BlossomConfig::from_profile( + BlossomProfile::simulator(["http://127.0.0.1:3000"]).unwrap(), + )) + .is_err() + ); + assert!(poisoned.profile_kind().is_none()); + } + #[cfg(feature = "radrootsd")] #[test] fn daemon_configuration_is_inert_explicit_and_redacted() { diff --git a/crates/sdk/tests/package_boundary.rs b/crates/sdk/tests/package_boundary.rs @@ -91,6 +91,7 @@ fn manifest_has_exact_feature_vocabulary_and_explicit_optional_activation() { for activation in [ "dep:radroots_storage_sqlite", "dep:radroots_sync", + "dep:radroots_blossom", "dep:radroots_nostr", "dep:radroots_transport_nostr", "dep:radroots_nostr_connect", @@ -98,6 +99,7 @@ fn manifest_has_exact_feature_vocabulary_and_explicit_optional_activation() { "dep:reqwest", "dep:serde", "dep:serde_json", + "dep:tokio", "dep:radroots_geonames", ] { assert!( @@ -152,6 +154,7 @@ fn package_contains_only_reachable_sources_and_registered_targets() { rust_files(&root.join("src")), BTreeSet::from([ "adapters/mod.rs".to_owned(), + "adapters/blossom.rs".to_owned(), "adapters/radrootsd.rs".to_owned(), "capability.rs".to_owned(), "client.rs".to_owned(), diff --git a/crates/sdk/tests/public_api.rs b/crates/sdk/tests/public_api.rs @@ -74,6 +74,19 @@ fn public_native_type_snapshot_uses_contextual_names() { .chain([ type_name::<radroots_sdk::signing::AuthorizationHeader>(), type_name::<radroots_sdk::signing::BlossomSigningError>(), + type_name::<radroots_sdk::transport::BlossomCancellation>(), + type_name::<radroots_sdk::transport::BlossomConfig>(), + type_name::<radroots_sdk::transport::BlossomEndpoint>(), + type_name::<radroots_sdk::transport::BlossomEndpointPolicy>(), + type_name::<radroots_sdk::transport::BlossomError>(), + type_name::<radroots_sdk::transport::BlossomErrorKind>(), + type_name::<radroots_sdk::transport::BlossomImageDimensions>(), + type_name::<radroots_sdk::transport::BlossomPhase>(), + type_name::<radroots_sdk::transport::BlossomProfile>(), + type_name::<radroots_sdk::transport::BlossomProfileKind>(), + type_name::<radroots_sdk::transport::BlossomSlot>(), + type_name::<radroots_sdk::transport::BlossomUploadReceipt>(), + type_name::<radroots_sdk::transport::BlossomUploadRequest>(), ]) .collect::<BTreeSet<_>>(); #[cfg(feature = "radrootsd")] @@ -153,7 +166,7 @@ const fn expected_public_type_count() -> usize { #[cfg(feature = "nostr")] let count = count + 1; #[cfg(feature = "blossom")] - let count = count + 2; + let count = count + 15; #[cfg(feature = "radrootsd")] let count = count + 5; count