commit 2dea7bb3a84417e1ed8f35bb285c7d2228104cbd
parent f7218237bb70e97fc9cd8f902dc8cffa9706d667
Author: triesap <tyson@radroots.org>
Date: Mon, 10 Aug 2026 18:03:02 +0000
feat: harden inbound Blossom retrieval
Verify canonical DNS, redirect, byte, media, and dimension bounds before accepting inbound images.
Persist exact artifacts atomically and revoke corrupt or obsolete cache receipts.
Diffstat:
13 files changed, 1592 insertions(+), 65 deletions(-)
diff --git a/crates/mobile_core/Cargo.toml b/crates/mobile_core/Cargo.toml
@@ -48,7 +48,7 @@ serde = { workspace = true, features = ["derive"] }
serde_json = { workspace = true }
sha2 = { workspace = true }
thiserror = { workspace = true }
-tokio = { workspace = true, optional = true, features = ["sync"] }
+tokio = { workspace = true, optional = true, features = ["fs", "io-util", "rt", "sync"] }
uuid = { workspace = true, features = ["v4"] }
url = { workspace = true }
diff --git a/crates/mobile_core/src/runtime/builder.rs b/crates/mobile_core/src/runtime/builder.rs
@@ -62,6 +62,8 @@ impl RuntimeBuilder {
}
self.store.validate_host_filesystem()?;
let options = self.store.sqlite_options()?;
+ #[cfg(feature = "mobile-social")]
+ let inbound_media_directory = self.store.owner_directory().join("inbound_media.v1");
let builder = radroots_sdk::ClientBuilder::sqlite(options)
.await
.map_err(RadrootsAppError::from_sdk)?;
@@ -69,6 +71,8 @@ impl RuntimeBuilder {
builder,
Some(self.store.public_key()),
#[cfg(feature = "mobile-social")]
+ Some(inbound_media_directory),
+ #[cfg(feature = "mobile-social")]
self.signer,
#[cfg(feature = "mobile-social")]
Some(self.relay_profile),
diff --git a/crates/mobile_core/src/runtime/mod.rs b/crates/mobile_core/src/runtime/mod.rs
@@ -8,6 +8,8 @@ pub mod store;
use chrono::Utc;
use radroots_identity::PublicKey;
use radroots_sdk::{Client, ClientBuilder};
+#[cfg(feature = "mobile-social")]
+use std::path::PathBuf;
use std::sync::{
RwLock,
atomic::{AtomicBool, Ordering},
@@ -27,12 +29,17 @@ pub struct RadrootsRuntime {
pub(crate) store_public_key: Option<PublicKey>,
#[cfg(feature = "mobile-social")]
pub(crate) settings_lock: tokio::sync::Mutex<()>,
+ #[cfg(feature = "mobile-social")]
+ pub(crate) inbound_media_directory: Option<PathBuf>,
+ #[cfg(feature = "mobile-social")]
+ pub(crate) inbound_media_lock: tokio::sync::Mutex<()>,
}
impl RadrootsRuntime {
pub(crate) fn from_client_builder(
builder: ClientBuilder,
store_public_key: Option<PublicKey>,
+ #[cfg(feature = "mobile-social")] inbound_media_directory: Option<PathBuf>,
#[cfg(feature = "mobile-social")] signer: Option<
std::sync::Arc<dyn radroots_signing::Signer>,
>,
@@ -77,6 +84,10 @@ impl RadrootsRuntime {
store_public_key,
#[cfg(feature = "mobile-social")]
settings_lock: tokio::sync::Mutex::new(()),
+ #[cfg(feature = "mobile-social")]
+ inbound_media_directory,
+ #[cfg(feature = "mobile-social")]
+ inbound_media_lock: tokio::sync::Mutex::new(()),
})
}
@@ -91,6 +102,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
@@ -28,6 +28,8 @@ pub use context::{
};
pub use cursor::{CursorError, CursorScope, TodayCursor, TodayCursorPosition};
pub use identity::{CARD_ID_SCHEMA_VERSION, CardId, CardIdError, CardSourceIdentity};
+#[cfg(feature = "mobile-social")]
+pub use media::Phase1LocalMediaArtifact;
pub use media::{
MediaReference, Phase1InboundMediaError, Phase1InboundMediaFailure, Phase1InboundMediaPending,
Phase1InboundMediaState, Phase1MediaArtifactId, Phase1MediaCacheIndex, Phase1MediaCachePolicy,
diff --git a/crates/mobile_core/src/runtime/product_surface/media.rs b/crates/mobile_core/src/runtime/product_surface/media.rs
@@ -6,11 +6,15 @@
//! an actual byte commitment and bound to the active retrieval configuration.
use std::collections::BTreeMap;
+#[cfg(feature = "mobile-social")]
+use std::path::{Path, PathBuf};
use radroots_blossom::{BlobUrl, MediaType, Sha256, descriptor::ByteCommitment};
use serde::{Deserialize, Serialize};
use sha2::{Digest, Sha256 as Sha256Hasher};
use thiserror::Error;
+#[cfg(feature = "mobile-social")]
+use tokio::io::AsyncWriteExt;
const MEDIA_REFERENCE_SCHEMA_VERSION: u16 = 1;
const MEDIA_RECEIPT_SCHEMA_VERSION: u16 = 1;
@@ -18,7 +22,11 @@ const MEDIA_CACHE_SCHEMA_VERSION: u16 = 1;
const MEDIA_URL_MAX_BYTES: usize = 8_192;
const MEDIA_ALT_MAX_BYTES: usize = 2_048;
const MEDIA_FAILURE_CODE_MAX_BYTES: usize = 96;
+const MEDIA_DIMENSION_MAX_EDGE: u32 = 16_384;
+const MEDIA_DIMENSION_MAX_PIXELS: u64 = 100_000_000;
const MEDIA_REFERENCE_FINGERPRINT_DOMAIN: &[u8] = b"radroots.inbound-media-reference.v1\0";
+#[cfg(feature = "mobile-social")]
+const MEDIA_CACHE_EXTENSIONS: &[&str] = &["gif", "jpg", "png", "webp"];
#[derive(Clone, Copy, Debug, Deserialize, Eq, Hash, Ord, PartialEq, PartialOrd, Serialize)]
#[serde(transparent)]
@@ -400,6 +408,9 @@ impl Phase1VerifiedMediaReceipt {
.ok_or(Phase1InboundMediaError::InvalidReference)?
.as_str()
.to_owned();
+ if canonical_extension(commitment.media_type()) != Some(extension.as_str()) {
+ return Err(Phase1InboundMediaError::MetadataMismatch);
+ }
if canonical_final_url.hash_path().hash() != commitment.sha256() {
return Err(Phase1InboundMediaError::MetadataMismatch);
}
@@ -480,6 +491,10 @@ impl Phase1VerifiedMediaReceipt {
.extension()
.is_none_or(|value| value.as_str() != self.extension)
|| !canonical_media_type(&self.media_type)
+ || MediaType::parse(&self.media_type)
+ .ok()
+ .and_then(|value| canonical_extension(&value))
+ != Some(self.extension.as_str())
|| validate_dimensions(Some(self.width), Some(self.height)).is_err()
{
return Err(Phase1InboundMediaError::CorruptReceipt);
@@ -522,6 +537,64 @@ impl Phase1VerifiedMediaReceipt {
pub const fn height(&self) -> u32 {
self.height
}
+
+ pub const fn verified_at_unix_ms(&self) -> u64 {
+ self.verified_at_unix_ms
+ }
+}
+
+/// One immutable exact-byte artifact in the authenticated user's local cache.
+#[cfg(feature = "mobile-social")]
+#[derive(Clone, Eq, PartialEq)]
+pub struct Phase1LocalMediaArtifact {
+ artifact_id: Phase1MediaArtifactId,
+ local_path: PathBuf,
+ byte_size: u64,
+ media_type: String,
+ width: u32,
+ height: u32,
+}
+
+#[cfg(feature = "mobile-social")]
+impl std::fmt::Debug for Phase1LocalMediaArtifact {
+ fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
+ formatter
+ .debug_struct("Phase1LocalMediaArtifact")
+ .field("artifact_id", &self.artifact_id)
+ .field("local_path", &"<redacted>")
+ .field("byte_size", &self.byte_size)
+ .field("media_type", &self.media_type)
+ .field("width", &self.width)
+ .field("height", &self.height)
+ .finish()
+ }
+}
+
+#[cfg(feature = "mobile-social")]
+impl Phase1LocalMediaArtifact {
+ pub const fn artifact_id(&self) -> Phase1MediaArtifactId {
+ self.artifact_id
+ }
+
+ pub fn local_path(&self) -> &Path {
+ self.local_path.as_path()
+ }
+
+ pub const fn byte_size(&self) -> u64 {
+ self.byte_size
+ }
+
+ pub fn media_type(&self) -> &str {
+ self.media_type.as_str()
+ }
+
+ pub const fn width(&self) -> u32 {
+ self.width
+ }
+
+ pub const fn height(&self) -> u32 {
+ self.height
+ }
}
#[derive(Clone, Debug, Default, Deserialize, Eq, PartialEq, Serialize)]
@@ -800,7 +873,7 @@ impl Phase1MediaCacheIndex {
{
return Err(Phase1InboundMediaError::ArtifactCollision);
}
- self.entries.insert(key, entry);
+ self.entries.insert(key.clone(), entry);
let mut evicted = Vec::new();
while self.entries.len() > policy.max_artifacts as usize
|| self.total_bytes()? > policy.max_bytes
@@ -808,6 +881,7 @@ impl Phase1MediaCacheIndex {
let oldest = self
.entries
.iter()
+ .filter(|(candidate, _)| candidate.as_str() != key)
.min_by_key(|(key, entry)| (entry.last_accessed_at_unix_ms, key.as_str()))
.map(|(key, _)| key.clone())
.ok_or(Phase1InboundMediaError::CorruptState)?;
@@ -969,6 +1043,190 @@ pub enum Phase1InboundMediaError {
CorruptState,
#[error("inbound media schema version is unsupported")]
UnsupportedSchema,
+ #[error("inbound media cache directory is unavailable")]
+ CacheUnavailable,
+ #[error("inbound media cache filesystem operation failed")]
+ CacheIo,
+ #[error("inbound media cache artifact is corrupt")]
+ CorruptArtifact,
+}
+
+#[cfg(feature = "mobile-social")]
+pub(crate) async fn write_verified_artifact(
+ directory: &Path,
+ receipt: &Phase1VerifiedMediaReceipt,
+ bytes: &[u8],
+) -> Result<Phase1LocalMediaArtifact, Phase1InboundMediaError> {
+ receipt.validate_intrinsic()?;
+ if bytes.len() as u64 != receipt.byte_size
+ || Sha256::digest(bytes).to_hex() != receipt.observed_sha256
+ {
+ return Err(Phase1InboundMediaError::CorruptArtifact);
+ }
+ ensure_cache_directory(directory).await?;
+ let final_path = artifact_path(directory, receipt.artifact_id, receipt.extension.as_str())?;
+ match tokio::fs::symlink_metadata(&final_path).await {
+ Ok(_) => {
+ verify_artifact_file(&final_path, receipt).await?;
+ return Ok(local_artifact(final_path, receipt));
+ }
+ Err(error) if error.kind() == std::io::ErrorKind::NotFound => {}
+ Err(_) => return Err(Phase1InboundMediaError::CacheIo),
+ }
+
+ let temporary_path = directory.join(format!(
+ ".{}.{}.tmp",
+ receipt.artifact_id.to_hex(),
+ uuid::Uuid::new_v4().simple()
+ ));
+ let mut temporary = tokio::fs::OpenOptions::new()
+ .create_new(true)
+ .write(true)
+ .open(&temporary_path)
+ .await
+ .map_err(|_| Phase1InboundMediaError::CacheIo)?;
+ let write_result = async {
+ temporary
+ .write_all(bytes)
+ .await
+ .map_err(|_| Phase1InboundMediaError::CacheIo)?;
+ temporary
+ .flush()
+ .await
+ .map_err(|_| Phase1InboundMediaError::CacheIo)?;
+ temporary
+ .sync_all()
+ .await
+ .map_err(|_| Phase1InboundMediaError::CacheIo)?;
+ drop(temporary);
+ match tokio::fs::hard_link(&temporary_path, &final_path).await {
+ Ok(()) => {}
+ Err(error) if error.kind() == std::io::ErrorKind::AlreadyExists => {
+ verify_artifact_file(&final_path, receipt).await?;
+ }
+ Err(_) => return Err(Phase1InboundMediaError::CacheIo),
+ }
+ tokio::fs::remove_file(&temporary_path)
+ .await
+ .map_err(|_| Phase1InboundMediaError::CacheIo)?;
+ sync_cache_directory(directory).await?;
+ verify_artifact_file(&final_path, receipt).await
+ }
+ .await;
+ if write_result.is_err() {
+ let _ = tokio::fs::remove_file(&temporary_path).await;
+ }
+ write_result?;
+ Ok(local_artifact(final_path, receipt))
+}
+
+#[cfg(feature = "mobile-social")]
+pub(crate) async fn remove_artifact_files(
+ directory: &Path,
+ artifact_id: Phase1MediaArtifactId,
+) -> Result<(), Phase1InboundMediaError> {
+ ensure_cache_directory(directory).await?;
+ for extension in MEDIA_CACHE_EXTENSIONS {
+ let path = artifact_path(directory, artifact_id, extension)?;
+ match tokio::fs::symlink_metadata(&path).await {
+ Ok(metadata) if metadata.file_type().is_symlink() || !metadata.is_file() => {
+ return Err(Phase1InboundMediaError::CorruptArtifact);
+ }
+ Ok(_) => tokio::fs::remove_file(path)
+ .await
+ .map_err(|_| Phase1InboundMediaError::CacheIo)?,
+ Err(error) if error.kind() == std::io::ErrorKind::NotFound => {}
+ Err(_) => return Err(Phase1InboundMediaError::CacheIo),
+ }
+ }
+ sync_cache_directory(directory).await
+}
+
+#[cfg(feature = "mobile-social")]
+pub(crate) async fn verified_artifact(
+ directory: &Path,
+ receipt: &Phase1VerifiedMediaReceipt,
+) -> Result<Phase1LocalMediaArtifact, Phase1InboundMediaError> {
+ receipt.validate_intrinsic()?;
+ ensure_cache_directory(directory).await?;
+ let path = artifact_path(directory, receipt.artifact_id, receipt.extension.as_str())?;
+ verify_artifact_file(&path, receipt).await?;
+ Ok(local_artifact(path, receipt))
+}
+
+#[cfg(feature = "mobile-social")]
+async fn ensure_cache_directory(directory: &Path) -> Result<(), Phase1InboundMediaError> {
+ match tokio::fs::create_dir(directory).await {
+ Ok(()) => {}
+ Err(error) if error.kind() == std::io::ErrorKind::AlreadyExists => {}
+ Err(_) => return Err(Phase1InboundMediaError::CacheIo),
+ }
+ let metadata = tokio::fs::symlink_metadata(directory)
+ .await
+ .map_err(|_| Phase1InboundMediaError::CacheIo)?;
+ if metadata.file_type().is_symlink() || !metadata.is_dir() {
+ return Err(Phase1InboundMediaError::CorruptArtifact);
+ }
+ Ok(())
+}
+
+#[cfg(feature = "mobile-social")]
+fn artifact_path(
+ directory: &Path,
+ artifact_id: Phase1MediaArtifactId,
+ extension: &str,
+) -> Result<PathBuf, Phase1InboundMediaError> {
+ if !MEDIA_CACHE_EXTENSIONS.contains(&extension) {
+ return Err(Phase1InboundMediaError::CorruptArtifact);
+ }
+ Ok(directory.join(format!("{}.{}", artifact_id.to_hex(), extension)))
+}
+
+#[cfg(feature = "mobile-social")]
+async fn verify_artifact_file(
+ path: &Path,
+ receipt: &Phase1VerifiedMediaReceipt,
+) -> Result<(), Phase1InboundMediaError> {
+ let metadata = tokio::fs::symlink_metadata(path)
+ .await
+ .map_err(|_| Phase1InboundMediaError::CorruptArtifact)?;
+ if metadata.file_type().is_symlink()
+ || !metadata.is_file()
+ || metadata.len() != receipt.byte_size
+ {
+ return Err(Phase1InboundMediaError::CorruptArtifact);
+ }
+ let bytes = tokio::fs::read(path)
+ .await
+ .map_err(|_| Phase1InboundMediaError::CacheIo)?;
+ if Sha256::digest(bytes.as_slice()).to_hex() != receipt.observed_sha256 {
+ return Err(Phase1InboundMediaError::CorruptArtifact);
+ }
+ Ok(())
+}
+
+#[cfg(feature = "mobile-social")]
+async fn sync_cache_directory(directory: &Path) -> Result<(), Phase1InboundMediaError> {
+ let directory = directory.to_path_buf();
+ tokio::task::spawn_blocking(move || std::fs::File::open(directory)?.sync_all())
+ .await
+ .map_err(|_| Phase1InboundMediaError::CacheIo)?
+ .map_err(|_| Phase1InboundMediaError::CacheIo)
+}
+
+#[cfg(feature = "mobile-social")]
+fn local_artifact(
+ local_path: PathBuf,
+ receipt: &Phase1VerifiedMediaReceipt,
+) -> Phase1LocalMediaArtifact {
+ Phase1LocalMediaArtifact {
+ artifact_id: receipt.artifact_id,
+ local_path,
+ byte_size: receipt.byte_size,
+ media_type: receipt.media_type.clone(),
+ width: receipt.width,
+ height: receipt.height,
+ }
}
fn validate_url_text(value: &str) -> Result<(), Phase1InboundMediaError> {
@@ -985,13 +1243,31 @@ fn canonical_media_type(value: &str) -> bool {
MediaType::parse(value).is_ok_and(|parsed| parsed.to_string() == value)
}
+fn canonical_extension(media_type: &MediaType) -> Option<&'static str> {
+ match media_type.as_str() {
+ "image/gif" => Some("gif"),
+ "image/jpeg" => Some("jpg"),
+ "image/png" => Some("png"),
+ "image/webp" => Some("webp"),
+ _ => None,
+ }
+}
+
fn validate_dimensions(
width: Option<u32>,
height: Option<u32>,
) -> Result<(), Phase1InboundMediaError> {
match (width, height) {
(None, None) => Ok(()),
- (Some(width), Some(height)) if width != 0 && height != 0 => Ok(()),
+ (Some(width), Some(height))
+ if width != 0
+ && height != 0
+ && width <= MEDIA_DIMENSION_MAX_EDGE
+ && height <= MEDIA_DIMENSION_MAX_EDGE
+ && u64::from(width) * u64::from(height) <= MEDIA_DIMENSION_MAX_PIXELS =>
+ {
+ Ok(())
+ }
_ => Err(Phase1InboundMediaError::InvalidDimensions),
}
}
@@ -1238,4 +1514,53 @@ mod tests {
let corrupt: Phase1MediaCacheIndex = serde_json::from_value(cache_value).unwrap();
assert_eq!(corrupt.status(), Err(Phase1InboundMediaError::CorruptState));
}
+
+ #[cfg(feature = "mobile-social")]
+ #[tokio::test]
+ async fn atomic_artifact_writes_converge_and_corruption_fails_closed() {
+ let bytes = b"GIF89a\x02\0\x03\0";
+ let hash = Sha256::digest(bytes).to_hex();
+ let structural = Phase1StructuralMediaReference::new(
+ format!("https://media.example/{hash}.gif"),
+ Some(hash.clone()),
+ Some("image/gif".to_owned()),
+ Some(2),
+ Some(3),
+ Some(bytes.len() as u64),
+ None,
+ )
+ .unwrap();
+ let receipt = Phase1VerifiedMediaReceipt::from_commitment(
+ &structural,
+ BlobUrl::parse(&format!("https://media.example/{hash}.gif")).unwrap(),
+ &ByteCommitment::from_bytes(bytes, MediaType::parse("image/gif").unwrap()),
+ 2,
+ 3,
+ configuration(4),
+ 1,
+ )
+ .unwrap();
+ let root = tempfile::tempdir().unwrap();
+ let directory = root.path().join("cache");
+ let (left, right) = tokio::join!(
+ write_verified_artifact(&directory, &receipt, bytes),
+ write_verified_artifact(&directory, &receipt, bytes),
+ );
+ let left = left.unwrap();
+ let right = right.unwrap();
+ assert_eq!(left, right);
+ assert_eq!(tokio::fs::read(left.local_path()).await.unwrap(), bytes);
+ assert_eq!(
+ std::fs::read_dir(&directory).unwrap().count(),
+ 1,
+ "no temporary file survives a converged write"
+ );
+ tokio::fs::write(left.local_path(), b"GIF89a\x03\0\x03\0")
+ .await
+ .unwrap();
+ assert_eq!(
+ verified_artifact(&directory, &receipt).await,
+ Err(Phase1InboundMediaError::CorruptArtifact)
+ );
+ }
}
diff --git a/crates/mobile_core/src/runtime/product_surface/outbox.rs b/crates/mobile_core/src/runtime/product_surface/outbox.rs
@@ -2732,6 +2732,7 @@ mod tests {
None,
None,
None,
+ None,
)
.unwrap()
}
@@ -2741,6 +2742,7 @@ mod tests {
ClientBuilder::memory_default(),
Some(PublicKey::from_hex(AUTHOR).unwrap()),
None,
+ None,
Some(profile),
None,
)
@@ -2755,6 +2757,7 @@ mod tests {
RadrootsRuntime::from_client_builder(
ClientBuilder::memory_default(),
Some(PublicKey::from_hex(AUTHOR).unwrap()),
+ None,
Some(std::sync::Arc::new(signer)),
None,
None,
@@ -3095,6 +3098,7 @@ mod tests {
let runtime = RadrootsRuntime::from_client_builder(
ClientBuilder::memory_default(),
Some(PublicKey::from_hex(AUTHOR).unwrap()),
+ None,
Some(std::sync::Arc::new(signer)),
Some(profile),
None,
diff --git a/crates/mobile_core/src/runtime/product_surface/today.rs b/crates/mobile_core/src/runtime/product_surface/today.rs
@@ -20,8 +20,14 @@ use sha2::{Digest, Sha256};
use thiserror::Error;
#[cfg(feature = "mobile-social")]
+use radroots_blossom::{BlobUrl, MediaType};
+#[cfg(feature = "mobile-social")]
use radroots_event::admission::ContractValidatedEvent;
#[cfg(feature = "mobile-social")]
+use radroots_sdk::transport::{
+ BlossomCancellation, BlossomError, BlossomImageDimensions, BlossomInboundRequest,
+};
+#[cfg(feature = "mobile-social")]
use radroots_sync::{
PullRequest,
ingest::{AdmissionDecision, AdmissionPolicy},
@@ -32,17 +38,20 @@ use radroots_transport::{
Target, outcome::FetchTargetState, source::FetchSelector, target::TargetSet,
};
+#[cfg(feature = "mobile-social")]
+use super::Phase1LocalMediaArtifact;
use super::{
CardId, CardLifecycleState, ClassifiedCard, CursorError, CursorScope, LocalAuthorOverlay,
LocalNetwork, LocalityEvidence, MeSnapshot, MediaReference, Phase1InboundMediaError,
Phase1InboundMediaFailure, Phase1InboundMediaPending, Phase1InboundMediaState,
- Phase1MediaArtifactId, Phase1MediaCacheIndex, Phase1MediaCachePolicy, Phase1MediaCacheStatus,
+ Phase1MediaArtifactId, Phase1MediaCacheIndex, Phase1MediaCacheStatus,
Phase1MediaConfigurationFingerprint, Phase1StructuralMediaReference,
- Phase1VerifiedMediaReceipt, ProductEventClassification, ProfileSummary, SearchResult,
- SearchResultType, SupportingProfile, ThreadEntry, ThreadReference, TimeRelevance, TodayCard,
- TodayCardType, TodayCursor, TodayCursorPosition, TodayPage, TodayRank, TodayRankInput,
- classify_admitted_event,
+ ProductEventClassification, ProfileSummary, SearchResult, SearchResultType, SupportingProfile,
+ ThreadEntry, ThreadReference, TimeRelevance, TodayCard, TodayCardType, TodayCursor,
+ TodayCursorPosition, TodayPage, TodayRank, TodayRankInput, classify_admitted_event,
};
+#[cfg(any(feature = "mobile-social", test))]
+use super::{Phase1MediaCachePolicy, Phase1VerifiedMediaReceipt};
use crate::runtime::RadrootsRuntime;
const TODAY_PROJECTION_ID: &str = "radroots.today.v1";
@@ -159,6 +168,9 @@ pub enum TodayError {
Serialization,
#[error(transparent)]
InboundMedia(#[from] Phase1InboundMediaError),
+ #[cfg(feature = "mobile-social")]
+ #[error(transparent)]
+ InboundRetrieval(#[from] BlossomError),
}
#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
@@ -627,7 +639,8 @@ impl RadrootsRuntime {
/// Atomically binds exact-byte evidence to the matching reference and
/// admits its content-addressed cache entry under the active LRU quota.
- pub async fn phase1_commit_media_receipt(
+ #[cfg(any(feature = "mobile-social", test))]
+ pub(crate) async fn phase1_commit_media_receipt(
&self,
context: &LocalNetwork,
reference_fingerprint: [u8; 32],
@@ -690,6 +703,8 @@ impl RadrootsRuntime {
context: &LocalNetwork,
artifact_id: Phase1MediaArtifactId,
) -> Result<bool, TodayError> {
+ #[cfg(feature = "mobile-social")]
+ let _guard = self.inbound_media_lock.lock().await;
let storage = self
.client
.storage()
@@ -703,6 +718,10 @@ impl RadrootsRuntime {
if cache_changed || references_changed {
persist_media_state(storage, context, generation, &mut state).await?;
}
+ #[cfg(feature = "mobile-social")]
+ if let Some(directory) = self.inbound_media_directory.as_deref() {
+ super::media::remove_artifact_files(directory, artifact_id).await?;
+ }
Ok(cache_changed || references_changed)
}
@@ -713,6 +732,8 @@ impl RadrootsRuntime {
context: &LocalNetwork,
configuration: Phase1MediaConfigurationFingerprint,
) -> Result<Vec<Phase1MediaArtifactId>, TodayError> {
+ #[cfg(feature = "mobile-social")]
+ let _guard = self.inbound_media_lock.lock().await;
let storage = self
.client
.storage()
@@ -729,6 +750,12 @@ impl RadrootsRuntime {
});
persist_media_state(storage, context, generation, &mut state).await?;
}
+ #[cfg(feature = "mobile-social")]
+ if let Some(directory) = self.inbound_media_directory.as_deref() {
+ for artifact_id in &removed {
+ super::media::remove_artifact_files(directory, *artifact_id).await?;
+ }
+ }
Ok(removed)
}
@@ -746,6 +773,198 @@ impl RadrootsRuntime {
state.media_cache.status().map_err(TodayError::from)
}
+ /// Resolves one renderable artifact only after rechecking the exact local
+ /// file. Missing or corrupt bytes atomically revoke all matching receipts.
+ #[cfg(feature = "mobile-social")]
+ pub async fn phase1_verified_media_artifact(
+ &self,
+ context: &LocalNetwork,
+ artifact_id: Phase1MediaArtifactId,
+ observed_at_unix_ms: u64,
+ ) -> Result<Option<Phase1LocalMediaArtifact>, TodayError> {
+ let _guard = self.inbound_media_lock.lock().await;
+ let directory = self
+ .inbound_media_directory
+ .as_deref()
+ .ok_or(Phase1InboundMediaError::CacheUnavailable)?;
+ let storage = self
+ .client
+ .storage()
+ .map_err(|_| TodayError::RuntimeUnavailable)?;
+ let generation = projection_generation()?;
+ let mut state = load_state(storage, context, generation)
+ .await?
+ .ok_or(TodayError::ProjectionMissing)?;
+ let Some(receipt) = verified_receipt(&state, artifact_id) else {
+ return Ok(None);
+ };
+ match super::media::verified_artifact(directory, &receipt).await {
+ Ok(artifact) => {
+ if state.media_cache.touch(artifact_id, observed_at_unix_ms)? {
+ persist_media_state(storage, context, generation, &mut state).await?;
+ }
+ Ok(Some(artifact))
+ }
+ Err(error) => {
+ state.media_cache.invalidate_artifact(artifact_id);
+ invalidate_artifact_references(&mut state, artifact_id);
+ persist_media_state(storage, context, generation, &mut state).await?;
+ let _ = super::media::remove_artifact_files(directory, artifact_id).await;
+ Err(error.into())
+ }
+ }
+ }
+
+ /// Completes one bounded BUD-01 retrieval, exact-byte verification, and
+ /// atomic content-addressed cache commit under the configured Blossom slot.
+ #[cfg(feature = "mobile-social")]
+ pub async fn phase1_retrieve_media(
+ &self,
+ context: &LocalNetwork,
+ reference_fingerprint: [u8; 32],
+ operation_id: [u8; 16],
+ policy: Phase1MediaCachePolicy,
+ cancellation: BlossomCancellation,
+ ) -> Result<Phase1LocalMediaArtifact, TodayError> {
+ let _guard = self.inbound_media_lock.lock().await;
+ let directory = self
+ .inbound_media_directory
+ .as_deref()
+ .ok_or(Phase1InboundMediaError::CacheUnavailable)?;
+ let blossom = self
+ .client
+ .blossom()
+ .map_err(|_| TodayError::RuntimeUnavailable)?
+ .cloned()
+ .ok_or(TodayError::RuntimeUnavailable)?;
+ let sdk_configuration = blossom
+ .config_fingerprint()
+ .ok_or(TodayError::RuntimeUnavailable)?;
+ let configuration =
+ Phase1MediaConfigurationFingerprint::new(*sdk_configuration.as_bytes())?;
+ let structural = load_structural_reference(self, context, reference_fingerprint).await?;
+ let started_at_unix_ms = inbound_now_unix_ms()?;
+ self.phase1_begin_media_retrieval(
+ context,
+ reference_fingerprint,
+ Phase1InboundMediaPending::new(operation_id, configuration, started_at_unix_ms)?,
+ )
+ .await?;
+ let request = match inbound_request(&structural) {
+ Ok(request) => request,
+ Err(error) => {
+ record_inbound_failure(
+ self,
+ context,
+ reference_fingerprint,
+ operation_id,
+ "invalid_reference",
+ false,
+ )
+ .await;
+ return Err(error);
+ }
+ };
+ let sdk_receipt = match blossom.retrieve(request, cancellation).await {
+ Ok(receipt) => receipt,
+ Err(error) => {
+ record_inbound_failure(
+ self,
+ context,
+ reference_fingerprint,
+ operation_id,
+ error.code().trim_start_matches("blossom_"),
+ error.retryable(),
+ )
+ .await;
+ return Err(error.into());
+ }
+ };
+ if sdk_receipt.config_fingerprint() != sdk_configuration {
+ record_inbound_failure(
+ self,
+ context,
+ reference_fingerprint,
+ operation_id,
+ "configuration_changed",
+ false,
+ )
+ .await;
+ return Err(Phase1InboundMediaError::ConfigurationMismatch.into());
+ }
+ let dimensions = sdk_receipt.dimensions();
+ let receipt = match Phase1VerifiedMediaReceipt::from_commitment(
+ &structural,
+ sdk_receipt.final_url().clone(),
+ sdk_receipt.commitment(),
+ dimensions.width(),
+ dimensions.height(),
+ configuration,
+ sdk_receipt.verified_at_unix_ms(),
+ ) {
+ Ok(receipt) => receipt,
+ Err(error) => {
+ record_inbound_failure(
+ self,
+ context,
+ reference_fingerprint,
+ operation_id,
+ "verification_failed",
+ false,
+ )
+ .await;
+ return Err(error.into());
+ }
+ };
+ let artifact =
+ match super::media::write_verified_artifact(directory, &receipt, sdk_receipt.bytes())
+ .await
+ {
+ Ok(artifact) => artifact,
+ Err(error) => {
+ record_inbound_failure(
+ self,
+ context,
+ reference_fingerprint,
+ operation_id,
+ "cache_write_failed",
+ true,
+ )
+ .await;
+ return Err(error.into());
+ }
+ };
+ let evicted = match self
+ .phase1_commit_media_receipt(
+ context,
+ reference_fingerprint,
+ operation_id,
+ receipt,
+ policy,
+ sdk_receipt.verified_at_unix_ms(),
+ )
+ .await
+ {
+ Ok(evicted) => evicted,
+ Err(error) => {
+ record_inbound_failure(
+ self,
+ context,
+ reference_fingerprint,
+ operation_id,
+ "cache_commit_failed",
+ true,
+ )
+ .await;
+ return Err(error);
+ }
+ };
+ for artifact_id in evicted {
+ super::media::remove_artifact_files(directory, artifact_id).await?;
+ }
+ Ok(artifact)
+ }
+
/// Persists active-author delivery state as a local-only Today overlay.
pub async fn phase1_set_local_author_overlay(
&self,
@@ -793,6 +1012,88 @@ impl RadrootsRuntime {
}
#[cfg(feature = "mobile-social")]
+async fn load_structural_reference(
+ runtime: &RadrootsRuntime,
+ context: &LocalNetwork,
+ reference_fingerprint: [u8; 32],
+) -> Result<Phase1StructuralMediaReference, TodayError> {
+ let storage = runtime
+ .client
+ .storage()
+ .map_err(|_| TodayError::RuntimeUnavailable)?;
+ let state = load_state(storage, context, projection_generation()?)
+ .await?
+ .ok_or(TodayError::ProjectionMissing)?;
+ state
+ .cards
+ .iter()
+ .flat_map(|projected| projected.card.media.iter())
+ .chain(
+ state
+ .profiles
+ .values()
+ .flat_map(|profile| [&profile.picture, &profile.banner].into_iter().flatten()),
+ )
+ .find(|media| media.structural().fingerprint() == &reference_fingerprint)
+ .map(|media| media.structural().clone())
+ .ok_or(TodayError::InvalidRequest)
+}
+
+#[cfg(feature = "mobile-social")]
+fn inbound_request(
+ structural: &Phase1StructuralMediaReference,
+) -> Result<BlossomInboundRequest, TodayError> {
+ let url = BlobUrl::parse(structural.source_url())
+ .map_err(|_| Phase1InboundMediaError::InvalidReference)?;
+ let media_type = structural
+ .expected_media_type()
+ .map(MediaType::parse)
+ .transpose()
+ .map_err(|_| Phase1InboundMediaError::InvalidMediaType)?;
+ let dimensions = match (structural.expected_width(), structural.expected_height()) {
+ (Some(width), Some(height)) => Some(BlossomImageDimensions::new(width, height)?),
+ (None, None) => None,
+ _ => return Err(Phase1InboundMediaError::InvalidDimensions.into()),
+ };
+ BlossomInboundRequest::new(url, media_type, structural.expected_byte_size(), dimensions)
+ .map_err(TodayError::from)
+}
+
+#[cfg(feature = "mobile-social")]
+fn inbound_now_unix_ms() -> Result<u64, TodayError> {
+ use std::time::{SystemTime, UNIX_EPOCH};
+
+ SystemTime::now()
+ .duration_since(UNIX_EPOCH)
+ .ok()
+ .and_then(|duration| u64::try_from(duration.as_millis()).ok())
+ .filter(|value| *value != 0)
+ .ok_or(TodayError::RuntimeUnavailable)
+}
+
+#[cfg(feature = "mobile-social")]
+async fn record_inbound_failure(
+ runtime: &RadrootsRuntime,
+ context: &LocalNetwork,
+ reference_fingerprint: [u8; 32],
+ operation_id: [u8; 16],
+ safe_code: &str,
+ retryable: bool,
+) {
+ let Ok(failed_at_unix_ms) = inbound_now_unix_ms() else {
+ return;
+ };
+ let Ok(failure) =
+ Phase1InboundMediaFailure::new(operation_id, safe_code, retryable, failed_at_unix_ms)
+ else {
+ return;
+ };
+ let _ = runtime
+ .phase1_fail_media_retrieval(context, reference_fingerprint, failure)
+ .await;
+}
+
+#[cfg(feature = "mobile-social")]
struct TodayAdmissionPolicy;
#[cfg(feature = "mobile-social")]
@@ -949,6 +1250,29 @@ fn invalidate_artifact_references(
changed
}
+#[cfg(feature = "mobile-social")]
+fn verified_receipt(
+ state: &TodayProjectionState,
+ artifact_id: Phase1MediaArtifactId,
+) -> Option<Phase1VerifiedMediaReceipt> {
+ state
+ .cards
+ .iter()
+ .flat_map(|projected| projected.card.media.iter())
+ .chain(
+ state
+ .profiles
+ .values()
+ .flat_map(|profile| [&profile.picture, &profile.banner].into_iter().flatten()),
+ )
+ .find_map(|media| match media.retrieval() {
+ Phase1InboundMediaState::Verified(receipt) if receipt.artifact_id() == artifact_id => {
+ Some((**receipt).clone())
+ }
+ _ => None,
+ })
+}
+
async fn persist_media_state(
storage: &dyn radroots_storage::Storage,
context: &LocalNetwork,
@@ -1691,6 +2015,10 @@ mod tests {
use super::*;
#[cfg(feature = "mobile-social")]
use std::sync::{Arc, Mutex, RwLock, atomic::AtomicBool};
+ #[cfg(feature = "mobile-social")]
+ use tokio::io::{AsyncReadExt, AsyncWriteExt};
+ #[cfg(feature = "mobile-social")]
+ use tokio::net::TcpListener;
use nostr::secp256k1::Message;
use nostr::{Keys, SECP256K1};
@@ -2026,6 +2354,42 @@ mod tests {
}
#[cfg(feature = "mobile-social")]
+ 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
+ }
+
+ #[cfg(feature = "mobile-social")]
+ async fn serve_one_blob(listener: TcpListener, bytes: Vec<u8>) {
+ let (mut stream, _) = listener.accept().await.expect("accept");
+ let mut request = Vec::new();
+ while !request.windows(4).any(|value| value == b"\r\n\r\n") {
+ let mut chunk = [0_u8; 1024];
+ let read = stream.read(&mut chunk).await.expect("request read");
+ assert_ne!(read, 0);
+ request.extend_from_slice(&chunk[..read]);
+ assert!(request.len() <= 64 * 1024);
+ }
+ let headers = String::from_utf8_lossy(&request);
+ assert!(headers.starts_with("GET /"));
+ assert!(
+ headers
+ .to_ascii_lowercase()
+ .contains("accept-encoding: identity")
+ );
+ assert!(!headers.to_ascii_lowercase().contains("authorization:"));
+ let response = 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(response.as_bytes()).await.expect("head");
+ stream.write_all(&bytes).await.expect("body");
+ stream.shutdown().await.expect("close");
+ }
+
+ #[cfg(feature = "mobile-social")]
#[tokio::test]
async fn relay_sync_fetches_the_exact_today_selector_and_projects_real_events() {
let source = Arc::new(TodaySource {
@@ -2044,6 +2408,8 @@ mod tests {
platform_app: RwLock::new(None),
store_public_key: None,
settings_lock: tokio::sync::Mutex::new(()),
+ inbound_media_directory: None,
+ inbound_media_lock: tokio::sync::Mutex::new(()),
};
let context = context(None, 1);
@@ -2753,6 +3119,165 @@ mod tests {
reopened.shutdown().await.expect("shutdown reopened");
}
+ #[cfg(feature = "mobile-social")]
+ #[tokio::test]
+ async fn hardened_retrieval_atomically_writes_and_invalidates_the_local_artifact() {
+ use crate::runtime::{
+ builder::RuntimeBuilder,
+ store::{MobileUserStoreConfig, ProtectedDataAvailability},
+ };
+ use radroots_sdk::transport::{
+ BlossomCancellation, BlossomConfig, BlossomEndpointAuthority, BlossomHostKind,
+ BlossomProfile,
+ };
+
+ let bytes = png(2, 3);
+ let hash = BlossomSha256::digest(bytes.as_slice()).to_hex();
+ let listener = TcpListener::bind("127.0.0.1:0").await.expect("bind");
+ let origin = format!("http://{}", listener.local_addr().expect("address"));
+ let url = format!("{origin}/{hash}.png");
+ let server = tokio::spawn(serve_one_blob(listener, bytes.clone()));
+ let root = tempfile::tempdir().expect("root");
+ let public_key = keys().public_key().to_string();
+ let store = MobileUserStoreConfig::from_encoded(
+ root.path(),
+ &public_key,
+ "6161616161616161616161616161616161616161616161616161616161616161",
+ 2_000_000_000_000,
+ ProtectedDataAvailability::Available,
+ )
+ .expect("store");
+ std::fs::create_dir_all(store.owner_directory()).expect("owner directory");
+ let owner_directory = store.owner_directory().to_path_buf();
+ let blossom = BlossomConfig::from_profile(
+ BlossomProfile::new(
+ BlossomHostKind::Simulator,
+ BlossomEndpointAuthority::LoopbackDevelopment,
+ origin,
+ std::iter::empty::<String>(),
+ )
+ .expect("profile"),
+ )
+ .with_network_policy(
+ std::time::Duration::from_millis(100),
+ std::time::Duration::from_millis(100),
+ 1,
+ std::time::Duration::from_millis(1),
+ )
+ .expect("network policy");
+ let runtime = RuntimeBuilder::new(store)
+ .blossom_config(blossom)
+ .build()
+ .await
+ .expect("runtime");
+ let context = context(None, 1);
+ ingest(
+ &runtime,
+ &context,
+ signed(
+ 0,
+ Vec::new(),
+ &format!(r#"{{"name":"retrieved","picture":"{url}"}}"#),
+ 2_000_000_000,
+ ),
+ 2_000_000_100,
+ )
+ .await;
+ let picture = runtime
+ .phase1_me(&context, &public_key, 2_000_000_200)
+ .await
+ .expect("me")
+ .profile
+ .expect("profile")
+ .picture
+ .expect("picture");
+ let artifact = runtime
+ .phase1_retrieve_media(
+ &context,
+ *picture.structural().fingerprint(),
+ [9; 16],
+ Phase1MediaCachePolicy::new(1_024, 8).unwrap(),
+ BlossomCancellation::default(),
+ )
+ .await
+ .expect("verified retrieval");
+ server.await.expect("server");
+ assert!(
+ artifact
+ .local_path()
+ .starts_with(owner_directory.join("inbound_media.v1"))
+ );
+ assert_eq!(tokio::fs::read(artifact.local_path()).await.unwrap(), bytes);
+ assert!(
+ !tokio::fs::symlink_metadata(artifact.local_path())
+ .await
+ .unwrap()
+ .file_type()
+ .is_symlink()
+ );
+ let verified = runtime
+ .phase1_me(&context, &public_key, 2_000_000_201)
+ .await
+ .expect("verified me")
+ .profile
+ .expect("verified profile")
+ .picture
+ .expect("verified picture");
+ assert!(matches!(
+ verified.retrieval(),
+ Phase1InboundMediaState::Verified(_)
+ ));
+ let mut corrupt = bytes.clone();
+ let last = corrupt.len() - 1;
+ corrupt[last] ^= 1;
+ tokio::fs::write(artifact.local_path(), corrupt)
+ .await
+ .expect("corrupt cached bytes");
+ assert!(matches!(
+ runtime
+ .phase1_verified_media_artifact(
+ &context,
+ artifact.artifact_id(),
+ 2_000_000_203_000,
+ )
+ .await,
+ Err(TodayError::InboundMedia(
+ Phase1InboundMediaError::CorruptArtifact
+ ))
+ ));
+ assert!(!artifact.local_path().exists());
+ let current_configuration = runtime
+ .phase1_media_cache_status(&context)
+ .await
+ .unwrap()
+ .configuration
+ .expect("configuration");
+ let replacement_bytes = if current_configuration.as_bytes() == &[1; 32] {
+ [2; 32]
+ } else {
+ [1; 32]
+ };
+ let replacement_configuration =
+ Phase1MediaConfigurationFingerprint::new(replacement_bytes).unwrap();
+ runtime
+ .phase1_invalidate_media_configuration(&context, replacement_configuration)
+ .await
+ .expect("invalidate configuration");
+ let unavailable = runtime
+ .phase1_me(&context, &public_key, 2_000_000_202)
+ .await
+ .expect("unavailable me")
+ .profile
+ .expect("unavailable profile")
+ .picture
+ .expect("unavailable picture");
+ assert!(matches!(
+ unavailable.retrieval(),
+ Phase1InboundMediaState::Unavailable
+ ));
+ runtime.shutdown().await.expect("shutdown");
+ }
+
#[tokio::test]
async fn fail_closed_requests_cursors_overlays_and_projection_guards_are_executable() {
let mut runtime = RadrootsRuntime::test_memory().expect("runtime");
diff --git a/crates/mobile_ffi/src/error.rs b/crates/mobile_ffi/src/error.rs
@@ -148,6 +148,13 @@ impl From<TodayError> for RadrootsAppError {
}
TodayError::RuntimeUnavailable => ("today_runtime_unavailable", true, &["retry"][..]),
TodayError::InboundMedia(_) => ("today_media_invalid", false, &["retry_media"][..]),
+ TodayError::InboundRetrieval(error) => {
+ if error.retryable() {
+ ("today_media_retrieval_failed", true, &["retry_media"][..])
+ } else {
+ ("today_media_retrieval_failed", false, &["review_media"][..])
+ }
+ }
TodayError::CorruptProjection | TodayError::Serialization | TodayError::Storage(_) => {
("today_state_failed", true, &["rebuild", "retry"][..])
}
diff --git a/crates/mobile_ffi/tests/local_mvp_real_io.rs b/crates/mobile_ffi/tests/local_mvp_real_io.rs
@@ -829,7 +829,7 @@ async fn prove_corrupted_media_fails(
assert_eq!(evidence.last_successful_state, "upload_verified");
assert_eq!(
evidence.error_code.as_deref(),
- Some("blossom_retrieved_bytes_mismatch")
+ Some("blossom_response_hash_mismatch")
);
assert_eq!(evidence.error_phase.as_deref(), Some("verification"));
assert!(evidence.possible_orphan);
diff --git a/crates/sdk/src/adapters/blossom.rs b/crates/sdk/src/adapters/blossom.rs
@@ -1,15 +1,18 @@
use std::{net::SocketAddr, time::Duration};
-use radroots_blossom::{BlobDescriptor, BlobUrl, MediaType};
+use radroots_blossom::{BlobDescriptor, BlobUrl, MediaType, Sha256, descriptor::ByteCommitment};
use reqwest::{
StatusCode,
- header::{ACCEPT, ACCEPT_ENCODING, AUTHORIZATION, CONTENT_LENGTH, CONTENT_TYPE, LOCATION},
+ header::{
+ ACCEPT, ACCEPT_ENCODING, AUTHORIZATION, CONTENT_ENCODING, CONTENT_LENGTH, CONTENT_TYPE,
+ LOCATION,
+ },
};
use crate::transport::{
BlossomCancellation, BlossomConfig, BlossomEndpoint, BlossomError, BlossomErrorKind,
- BlossomImageDimensions, BlossomPhase, BlossomUploadReceipt, BlossomUploadRequest,
- BlossomUploadTransaction,
+ BlossomImageDimensions, BlossomInboundReceipt, BlossomInboundRequest, BlossomPhase,
+ BlossomUploadReceipt, BlossomUploadRequest, BlossomUploadTransaction,
};
const MAX_RESOLVED_ADDRESSES: usize = 32;
@@ -245,7 +248,14 @@ async fn upload_bound_with_authorization(
{
Ok(descriptor) => break descriptor,
Err(error) if error.retryable() && upload_attempts < config.max_attempts() => {
- retry_delay(&config, upload_attempts, &cancellation, true).await?;
+ retry_delay(
+ &config,
+ upload_attempts,
+ &cancellation,
+ BlossomPhase::Upload,
+ true,
+ )
+ .await?;
}
Err(error) => return Err(error),
}
@@ -264,15 +274,26 @@ async fn upload_bound_with_authorization(
match retrieve_once(
&config,
verified_upload.url().as_blob_url().clone(),
- &request,
+ request.sha256(),
+ Some(request.media_type()),
+ Some(request.byte_size()),
+ Some(request.dimensions()),
&cancellation,
upload_attempts.saturating_add(retrieval_attempts),
+ true,
)
.await
{
- Ok(bytes) => break bytes,
+ Ok(retrieved) => break retrieved.bytes,
Err(error) if error.retryable() && retrieval_attempts < config.max_attempts() => {
- retry_delay(&config, retrieval_attempts, &cancellation, true).await?;
+ retry_delay(
+ &config,
+ retrieval_attempts,
+ &cancellation,
+ BlossomPhase::Retrieval,
+ true,
+ )
+ .await?;
}
Err(error) => return Err(error),
}
@@ -315,6 +336,67 @@ async fn upload_bound_with_authorization(
))
}
+pub(crate) async fn retrieve(
+ config: BlossomConfig,
+ request: BlossomInboundRequest,
+ cancellation: BlossomCancellation,
+) -> Result<BlossomInboundReceipt, BlossomError> {
+ if request
+ .expected_byte_size()
+ .is_some_and(|size| size > config.max_blob_bytes())
+ {
+ return Err(failure(
+ BlossomErrorKind::ResponseTooLarge,
+ BlossomPhase::Verification,
+ false,
+ false,
+ 0,
+ ));
+ }
+ let fingerprint = config.fingerprint();
+ let mut attempts = 0_u8;
+ let retrieved = loop {
+ ensure_not_cancelled(&cancellation, BlossomPhase::Retrieval, attempts, false)?;
+ attempts = attempts.saturating_add(1);
+ match retrieve_once(
+ &config,
+ request.url().clone(),
+ request.url().hash_path().hash(),
+ request.expected_media_type(),
+ request.expected_byte_size(),
+ request.expected_dimensions(),
+ &cancellation,
+ attempts,
+ false,
+ )
+ .await
+ {
+ Ok(retrieved) => break retrieved,
+ Err(error) if error.retryable() && attempts < config.max_attempts() => {
+ retry_delay(
+ &config,
+ attempts,
+ &cancellation,
+ BlossomPhase::Retrieval,
+ false,
+ )
+ .await?;
+ }
+ Err(error) => return Err(error),
+ }
+ };
+ let commitment = ByteCommitment::from_bytes(retrieved.bytes.as_slice(), retrieved.media_type);
+ Ok(BlossomInboundReceipt::new(
+ retrieved.final_url,
+ commitment,
+ retrieved.dimensions,
+ std::sync::Arc::from(retrieved.bytes),
+ fingerprint,
+ attempts,
+ crate::transport::blossom_now_unix_ms(),
+ ))
+}
+
#[cfg(test)]
async fn upload_with_authorization(
config: BlossomConfig,
@@ -474,20 +556,33 @@ fn verify_descriptor(
}
#[cfg_attr(coverage_nightly, coverage(off))]
+struct RetrievedBytes {
+ final_url: BlobUrl,
+ bytes: Vec<u8>,
+ media_type: MediaType,
+ dimensions: BlossomImageDimensions,
+}
+
+#[cfg_attr(coverage_nightly, coverage(off))]
+#[allow(clippy::too_many_arguments)]
async fn retrieve_once(
config: &BlossomConfig,
mut url: BlobUrl,
- request: &BlossomUploadRequest,
+ expected_sha256: Sha256,
+ expected_media_type: Option<&MediaType>,
+ expected_byte_size: Option<u64>,
+ expected_dimensions: Option<BlossomImageDimensions>,
cancellation: &BlossomCancellation,
attempt: u8,
-) -> Result<Vec<u8>, BlossomError> {
+ possible_orphan: bool,
+) -> Result<RetrievedBytes, 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,
+ possible_orphan,
attempt,
)
})?;
@@ -497,12 +592,15 @@ async fn retrieve_once(
cancellation,
BlossomPhase::Retrieval,
attempt,
- true,
+ possible_orphan,
)
.await?;
let pending = client
.get(url.as_str())
- .header(ACCEPT, request.media_type().as_str())
+ .header(
+ ACCEPT,
+ expected_media_type.map_or("image/*", MediaType::as_str),
+ )
.header(ACCEPT_ENCODING, "identity")
.send();
let response = tokio::select! {
@@ -512,11 +610,11 @@ async fn retrieve_once(
BlossomErrorKind::Cancelled,
BlossomPhase::Retrieval,
true,
- true,
+ possible_orphan,
attempt,
));
}
- response = pending => response.map_err(|error| request_error(error, BlossomPhase::Retrieval, true, attempt))?,
+ response = pending => response.map_err(|error| request_error(error, BlossomPhase::Retrieval, possible_orphan, attempt))?,
};
if response.status().is_redirection() {
if redirects == config.max_redirects() {
@@ -524,7 +622,7 @@ async fn retrieve_once(
BlossomErrorKind::RedirectLimit,
BlossomPhase::Retrieval,
false,
- true,
+ possible_orphan,
attempt,
));
}
@@ -537,7 +635,7 @@ async fn retrieve_once(
BlossomErrorKind::UnsafeRedirect,
BlossomPhase::Retrieval,
false,
- true,
+ possible_orphan,
attempt,
)
})?;
@@ -546,7 +644,7 @@ async fn retrieve_once(
BlossomErrorKind::UnsafeRedirect,
BlossomPhase::Retrieval,
false,
- true,
+ possible_orphan,
attempt,
)
})?;
@@ -555,7 +653,7 @@ async fn retrieve_once(
BlossomErrorKind::UnsafeRedirect,
BlossomPhase::Retrieval,
false,
- true,
+ possible_orphan,
attempt,
)
})?;
@@ -564,11 +662,11 @@ async fn retrieve_once(
BlossomErrorKind::UnsafeRedirect,
BlossomPhase::Retrieval,
false,
- true,
+ possible_orphan,
attempt,
)
})?;
- if next.hash_path().hash() != request.sha256()
+ if next.hash_path().hash() != expected_sha256
|| next.clone().approve().is_err()
|| config.profile().endpoint_for_blob(&next).is_none()
{
@@ -576,7 +674,7 @@ async fn retrieve_once(
BlossomErrorKind::UnsafeRedirect,
BlossomPhase::Retrieval,
false,
- true,
+ possible_orphan,
attempt,
));
}
@@ -587,7 +685,7 @@ async fn retrieve_once(
return Err(http_status_response_error(
response,
BlossomPhase::Retrieval,
- true,
+ possible_orphan,
attempt,
cancellation,
)
@@ -603,49 +701,144 @@ async fn retrieve_once(
BlossomErrorKind::MediaTypeMismatch,
BlossomPhase::Verification,
false,
- true,
+ possible_orphan,
attempt,
)
})?;
- if &actual_media_type != request.media_type() {
+ if expected_media_type.is_some_and(|expected| expected != &actual_media_type) {
return Err(failure(
BlossomErrorKind::MediaTypeMismatch,
BlossomPhase::Verification,
false,
- true,
+ possible_orphan,
+ attempt,
+ ));
+ }
+ let canonical_extension =
+ crate::transport::canonical_image_extension(&actual_media_type)
+ .map_err(|error| with_operation(error, possible_orphan, attempt))?;
+ if url
+ .hash_path()
+ .extension()
+ .is_none_or(|value| value.as_str() != canonical_extension)
+ {
+ return Err(failure(
+ BlossomErrorKind::MediaTypeMismatch,
+ BlossomPhase::Verification,
+ false,
+ possible_orphan,
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())
+ .get(CONTENT_ENCODING)
+ .is_some_and(|value| {
+ value
+ .to_str()
+ .map_or(true, |encoding| !encoding.eq_ignore_ascii_case("identity"))
+ })
{
return Err(failure(
+ BlossomErrorKind::ContentEncodingDenied,
+ BlossomPhase::Retrieval,
+ false,
+ possible_orphan,
+ attempt,
+ ));
+ }
+ let content_length = match response.headers().get(CONTENT_LENGTH) {
+ Some(value) => Some(
+ value
+ .to_str()
+ .ok()
+ .and_then(|value| value.parse::<u64>().ok())
+ .ok_or_else(|| {
+ failure(
+ BlossomErrorKind::ResponseSizeMismatch,
+ BlossomPhase::Retrieval,
+ false,
+ possible_orphan,
+ attempt,
+ )
+ })?,
+ ),
+ None => None,
+ };
+ if content_length.is_some_and(|size| size > config.max_blob_bytes()) {
+ return Err(failure(
BlossomErrorKind::ResponseTooLarge,
BlossomPhase::Retrieval,
false,
- true,
+ possible_orphan,
+ attempt,
+ ));
+ }
+ if expected_byte_size
+ .zip(content_length)
+ .is_some_and(|(expected, actual)| expected != actual)
+ {
+ return Err(failure(
+ BlossomErrorKind::ResponseSizeMismatch,
+ BlossomPhase::Verification,
+ false,
+ possible_orphan,
attempt,
));
}
- return read_bounded(
+ let bytes = read_bounded(
response,
usize::try_from(config.max_blob_bytes()).unwrap_or(usize::MAX),
cancellation,
BlossomPhase::Retrieval,
- true,
+ possible_orphan,
attempt,
)
- .await;
+ .await?;
+ let byte_size = u64::try_from(bytes.len()).unwrap_or(u64::MAX);
+ if expected_byte_size.is_some_and(|expected| expected != byte_size)
+ || content_length.is_some_and(|expected| expected != byte_size)
+ {
+ return Err(failure(
+ BlossomErrorKind::ResponseSizeMismatch,
+ BlossomPhase::Verification,
+ false,
+ possible_orphan,
+ attempt,
+ ));
+ }
+ if Sha256::digest(bytes.as_slice()) != expected_sha256 {
+ return Err(failure(
+ BlossomErrorKind::ResponseHashMismatch,
+ BlossomPhase::Verification,
+ false,
+ possible_orphan,
+ attempt,
+ ));
+ }
+ let dimensions = inspect_image(bytes.as_slice(), &actual_media_type)
+ .map_err(|error| with_operation(error, possible_orphan, attempt))?;
+ if expected_dimensions.is_some_and(|expected| expected != dimensions) {
+ return Err(failure(
+ BlossomErrorKind::DimensionMismatch,
+ BlossomPhase::Verification,
+ false,
+ possible_orphan,
+ attempt,
+ ));
+ }
+ return Ok(RetrievedBytes {
+ final_url: url,
+ bytes,
+ media_type: actual_media_type,
+ dimensions,
+ });
}
Err(failure(
BlossomErrorKind::RedirectLimit,
BlossomPhase::Retrieval,
false,
- true,
+ possible_orphan,
attempt,
))
}
@@ -836,6 +1029,7 @@ async fn retry_delay(
config: &BlossomConfig,
attempt: u8,
cancellation: &BlossomCancellation,
+ phase: BlossomPhase,
possible_orphan: bool,
) -> Result<(), BlossomError> {
let exponent = u32::from(attempt.saturating_sub(1)).min(16);
@@ -848,7 +1042,7 @@ async fn retry_delay(
biased;
_ = cancellation.cancelled() => Err(failure(
BlossomErrorKind::Cancelled,
- BlossomPhase::Upload,
+ phase,
true,
possible_orphan,
attempt,
@@ -1001,6 +1195,23 @@ pub(crate) fn verify_image(
media_type: &MediaType,
expected: BlossomImageDimensions,
) -> Result<(), BlossomError> {
+ let dimensions = inspect_image(bytes, media_type)?;
+ if dimensions != expected {
+ return Err(failure(
+ BlossomErrorKind::DimensionMismatch,
+ BlossomPhase::Verification,
+ false,
+ false,
+ 0,
+ ));
+ }
+ Ok(())
+}
+
+fn inspect_image(
+ bytes: &[u8],
+ media_type: &MediaType,
+) -> Result<BlossomImageDimensions, BlossomError> {
let detected = detect_image(bytes).ok_or_else(|| {
failure(
BlossomErrorKind::InvalidImageBytes,
@@ -1034,16 +1245,7 @@ pub(crate) fn verify_image(
0,
));
}
- if detected.1 != expected {
- return Err(failure(
- BlossomErrorKind::DimensionMismatch,
- BlossomPhase::Verification,
- false,
- false,
- 0,
- ));
- }
- Ok(())
+ Ok(detected.1)
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
@@ -1425,10 +1627,14 @@ mod tests {
)
.unwrap();
let delay = BlossomCancellation::default();
- assert!(retry_delay(&config, 1, &delay, false).await.is_ok());
+ assert!(
+ retry_delay(&config, 1, &delay, BlossomPhase::Upload, false)
+ .await
+ .is_ok()
+ );
delay.cancel();
assert_eq!(
- retry_delay(&config, 20, &delay, true)
+ retry_delay(&config, 20, &delay, BlossomPhase::Retrieval, true,)
.await
.expect_err("cancelled delay")
.kind(),
@@ -1823,6 +2029,215 @@ mod tests {
.expect("profile")
}
+ async fn spawn_inbound_server(
+ bytes: Vec<u8>,
+ content_type: &'static str,
+ content_encoding: Option<&'static str>,
+ declared_length: Option<usize>,
+ ) -> (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 task = tokio::spawn(async move {
+ let (mut stream, _) = listener.accept().await.expect("accept");
+ let request = read_request(&mut stream).await;
+ let encoding = content_encoding
+ .map(|value| format!("Content-Encoding: {value}\r\n"))
+ .unwrap_or_default();
+ let response = format!(
+ "HTTP/1.1 200 OK\r\nContent-Type: {content_type}\r\n{encoding}Content-Length: {}\r\nConnection: close\r\n\r\n",
+ declared_length.unwrap_or(bytes.len())
+ );
+ stream.write_all(response.as_bytes()).await.expect("head");
+ stream.write_all(&bytes).await.expect("body");
+ stream.shutdown().await.expect("close");
+ request
+ });
+ (origin, task)
+ }
+
+ fn inbound_request(
+ origin: &str,
+ bytes: &[u8],
+ dimensions: BlossomImageDimensions,
+ ) -> BlossomInboundRequest {
+ let hash = Sha256::digest(bytes);
+ BlossomInboundRequest::new(
+ BlobUrl::parse(format!("{origin}/{hash}.png").as_str()).expect("URL"),
+ Some(MediaType::parse("image/png").expect("media type")),
+ Some(bytes.len() as u64),
+ Some(dimensions),
+ )
+ .expect("inbound request")
+ }
+
+ #[tokio::test]
+ async fn inbound_retrieval_returns_only_exact_verified_image_bytes() {
+ let bytes = png(2, 3);
+ let (origin, server) = spawn_inbound_server(bytes.clone(), "image/png", None, None).await;
+ let slot = crate::transport::BlossomSlot::new();
+ slot.configure(config(origin.as_str())).expect("configure");
+ let receipt = slot
+ .retrieve(
+ inbound_request(
+ origin.as_str(),
+ bytes.as_slice(),
+ BlossomImageDimensions::new(2, 3).unwrap(),
+ ),
+ BlossomCancellation::default(),
+ )
+ .await
+ .expect("inbound receipt");
+ assert_eq!(receipt.bytes(), bytes.as_slice());
+ assert_eq!(receipt.commitment().sha256(), Sha256::digest(&bytes));
+ assert_eq!(receipt.commitment().size(), bytes.len() as u64);
+ assert_eq!(receipt.commitment().media_type().as_str(), "image/png");
+ assert_eq!(
+ receipt.dimensions(),
+ BlossomImageDimensions::new(2, 3).unwrap()
+ );
+ assert_eq!(receipt.attempts(), 1);
+ assert_eq!(
+ receipt.config_fingerprint(),
+ slot.config_fingerprint().expect("fingerprint")
+ );
+ let request = server.await.expect("server");
+ let headers = String::from_utf8_lossy(&request);
+ assert!(headers.starts_with("GET /"));
+ assert!(
+ headers
+ .to_ascii_lowercase()
+ .contains("accept-encoding: identity")
+ );
+ assert!(!headers.to_ascii_lowercase().contains("authorization:"));
+ }
+
+ #[tokio::test]
+ async fn inbound_retrieval_rejects_encoding_truncation_hash_and_dimensions() {
+ let exact = png(2, 3);
+ let cases = [
+ (
+ exact.clone(),
+ Some("gzip"),
+ None,
+ BlossomImageDimensions::new(2, 3).unwrap(),
+ BlossomErrorKind::ContentEncodingDenied,
+ ),
+ (
+ exact[..exact.len() - 1].to_vec(),
+ None,
+ Some(exact.len()),
+ BlossomImageDimensions::new(2, 3).unwrap(),
+ BlossomErrorKind::Transport,
+ ),
+ (
+ {
+ let mut altered = exact.clone();
+ let last = altered.len() - 1;
+ altered[last] ^= 1;
+ altered
+ },
+ None,
+ None,
+ BlossomImageDimensions::new(2, 3).unwrap(),
+ BlossomErrorKind::ResponseHashMismatch,
+ ),
+ (
+ exact.clone(),
+ None,
+ None,
+ BlossomImageDimensions::new(3, 2).unwrap(),
+ BlossomErrorKind::DimensionMismatch,
+ ),
+ ];
+ for (served, encoding, length, dimensions, expected) in cases {
+ let (origin, server) =
+ spawn_inbound_server(served, "image/png", encoding, length).await;
+ let slot = crate::transport::BlossomSlot::new();
+ slot.configure(config(origin.as_str())).expect("configure");
+ let error = slot
+ .retrieve(
+ inbound_request(origin.as_str(), exact.as_slice(), dimensions),
+ BlossomCancellation::default(),
+ )
+ .await
+ .expect_err("hostile response");
+ assert_eq!(error.kind(), expected);
+ assert!(!error.possible_orphan());
+ server.await.expect("server");
+ }
+ }
+
+ #[tokio::test]
+ async fn inbound_retry_and_cancellation_are_bounded_and_operation_scoped() {
+ let bytes = png(2, 3);
+ let listener = TcpListener::bind("127.0.0.1:0").await.expect("bind");
+ let origin = format!("http://{}", listener.local_addr().expect("address"));
+ let server_bytes = bytes.clone();
+ let server = tokio::spawn(async move {
+ for attempt in 0..2 {
+ let (mut stream, _) = listener.accept().await.expect("accept");
+ let _ = read_request(&mut stream).await;
+ if attempt == 0 {
+ stream
+ .write_all(b"HTTP/1.1 503 Service Unavailable\r\nContent-Length: 0\r\nConnection: close\r\n\r\n")
+ .await
+ .expect("retry response");
+ } else {
+ let response = format!(
+ "HTTP/1.1 200 OK\r\nContent-Type: image/png\r\nContent-Length: {}\r\nConnection: close\r\n\r\n",
+ server_bytes.len()
+ );
+ stream.write_all(response.as_bytes()).await.expect("head");
+ stream.write_all(&server_bytes).await.expect("body");
+ }
+ stream.shutdown().await.expect("close");
+ }
+ });
+ let slot = crate::transport::BlossomSlot::new();
+ slot.configure(
+ BlossomConfig::from_profile(simulator_profile(origin.as_str()))
+ .with_network_policy(
+ Duration::from_millis(100),
+ Duration::from_millis(100),
+ 2,
+ Duration::from_millis(1),
+ )
+ .unwrap(),
+ )
+ .unwrap();
+ let receipt = slot
+ .retrieve(
+ inbound_request(
+ origin.as_str(),
+ bytes.as_slice(),
+ BlossomImageDimensions::new(2, 3).unwrap(),
+ ),
+ BlossomCancellation::default(),
+ )
+ .await
+ .expect("retried retrieval");
+ assert_eq!(receipt.attempts(), 2);
+ server.await.expect("server");
+
+ let cancellation = BlossomCancellation::default();
+ cancellation.cancel();
+ let error = slot
+ .retrieve(
+ inbound_request(
+ origin.as_str(),
+ bytes.as_slice(),
+ BlossomImageDimensions::new(2, 3).unwrap(),
+ ),
+ cancellation,
+ )
+ .await
+ .expect_err("cancelled retrieval");
+ assert_eq!(error.kind(), BlossomErrorKind::Cancelled);
+ assert_eq!(error.phase(), BlossomPhase::Retrieval);
+ assert!(!error.possible_orphan());
+ }
+
#[tokio::test]
async fn non_mutating_probe_records_only_dns_transport_and_http_evidence() {
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
@@ -1991,7 +2406,7 @@ mod tests {
(
RetrievalResponse::Altered,
false,
- BlossomErrorKind::RetrievedBytesMismatch,
+ BlossomErrorKind::ResponseHashMismatch,
),
(
RetrievalResponse::WrongMime,
@@ -2001,7 +2416,7 @@ mod tests {
(
RetrievalResponse::Oversize,
false,
- BlossomErrorKind::ResponseTooLarge,
+ BlossomErrorKind::ResponseSizeMismatch,
),
(RetrievalResponse::Stall, false, BlossomErrorKind::Timeout),
];
diff --git a/crates/sdk/src/client.rs b/crates/sdk/src/client.rs
@@ -693,7 +693,14 @@ mod tests {
fn nostr_capabilities_require_directional_profile_authority_and_evidence() {
let slot = crate::transport::NostrSlot::new();
slot.configure(
- crate::transport::RelayProfile::public(Vec::<String>::new()).expect("public profile"),
+ crate::transport::RelayProfile::explicit(
+ crate::transport::RelayProfileKind::Public,
+ [(
+ radroots_transport_nostr::DEFAULT_PUBLIC_RELAY,
+ crate::transport::RelayAccess::ReadOnly,
+ )],
+ )
+ .expect("read-only public profile"),
)
.expect("configure slot");
let client = ClientBuilder::memory(generation())
diff --git a/crates/sdk/src/transport.rs b/crates/sdk/src/transport.rs
@@ -49,6 +49,10 @@ const MAX_BLOSSOM_ATTEMPTS: u8 = 5;
const MAX_BLOSSOM_TIMEOUT: Duration = Duration::from_secs(120);
#[cfg(feature = "blossom")]
const MAX_BLOSSOM_RETRY_DELAY: Duration = Duration::from_secs(30);
+#[cfg(feature = "blossom")]
+const MAX_BLOSSOM_IMAGE_EDGE: u32 = 16_384;
+#[cfg(feature = "blossom")]
+const MAX_BLOSSOM_IMAGE_PIXELS: u64 = 100_000_000;
/// Host environment executing Blossom operations.
#[cfg(feature = "blossom")]
@@ -426,6 +430,11 @@ pub struct BlossomConfigFingerprint(Sha256);
#[cfg(feature = "blossom")]
impl BlossomConfigFingerprint {
#[must_use]
+ pub const fn as_bytes(&self) -> &[u8; 32] {
+ self.0.as_bytes()
+ }
+
+ #[must_use]
pub fn to_hex(self) -> String {
self.0.to_hex()
}
@@ -445,7 +454,7 @@ fn append_fingerprint_field(material: &mut Vec<u8>, value: &[u8]) {
}
#[cfg(feature = "blossom")]
-fn blossom_now_unix_ms() -> u64 {
+pub(crate) fn blossom_now_unix_ms() -> u64 {
SystemTime::now()
.duration_since(UNIX_EPOCH)
.ok()
@@ -465,7 +474,12 @@ pub struct BlossomImageDimensions {
#[cfg(feature = "blossom")]
impl BlossomImageDimensions {
pub const fn new(width: u32, height: u32) -> Result<Self, BlossomError> {
- if width == 0 || height == 0 {
+ if width == 0
+ || height == 0
+ || width > MAX_BLOSSOM_IMAGE_EDGE
+ || height > MAX_BLOSSOM_IMAGE_EDGE
+ || width as u64 * height as u64 > MAX_BLOSSOM_IMAGE_PIXELS
+ {
return Err(BlossomError::configuration(
BlossomErrorKind::InvalidDimensions,
));
@@ -547,6 +561,73 @@ impl BlossomUploadRequest {
}
}
+/// Expected signed metadata for one bounded BUD-01 image retrieval.
+#[cfg(feature = "blossom")]
+#[derive(Clone, Debug, Eq, PartialEq)]
+pub struct BlossomInboundRequest {
+ url: BlobUrl,
+ expected_media_type: Option<MediaType>,
+ expected_byte_size: Option<u64>,
+ expected_dimensions: Option<BlossomImageDimensions>,
+}
+
+#[cfg(feature = "blossom")]
+impl BlossomInboundRequest {
+ pub fn new(
+ url: BlobUrl,
+ expected_media_type: Option<MediaType>,
+ expected_byte_size: Option<u64>,
+ expected_dimensions: Option<BlossomImageDimensions>,
+ ) -> Result<Self, BlossomError> {
+ url.clone()
+ .approve()
+ .map_err(|_| BlossomError::configuration(BlossomErrorKind::InvalidRequest))?;
+ if expected_byte_size == Some(0) {
+ return Err(BlossomError::configuration(
+ BlossomErrorKind::InvalidRequest,
+ ));
+ }
+ if let Some(media_type) = &expected_media_type {
+ let extension = canonical_image_extension(media_type)?;
+ if url
+ .hash_path()
+ .extension()
+ .is_none_or(|value| value.as_str() != extension)
+ {
+ return Err(BlossomError::configuration(
+ BlossomErrorKind::MediaTypeMismatch,
+ ));
+ }
+ }
+ Ok(Self {
+ url,
+ expected_media_type,
+ expected_byte_size,
+ expected_dimensions,
+ })
+ }
+
+ #[must_use]
+ pub const fn url(&self) -> &BlobUrl {
+ &self.url
+ }
+
+ #[must_use]
+ pub const fn expected_media_type(&self) -> Option<&MediaType> {
+ self.expected_media_type.as_ref()
+ }
+
+ #[must_use]
+ pub const fn expected_byte_size(&self) -> Option<u64> {
+ self.expected_byte_size
+ }
+
+ #[must_use]
+ pub const fn expected_dimensions(&self) -> Option<BlossomImageDimensions> {
+ self.expected_dimensions
+ }
+}
+
/// Immutable upload plan binding exact bytes to one complete configuration.
#[cfg(feature = "blossom")]
#[derive(Clone)]
@@ -872,7 +953,10 @@ pub enum BlossomErrorKind {
HttpStatus,
UnsafeRedirect,
RedirectLimit,
+ ContentEncodingDenied,
ResponseTooLarge,
+ ResponseSizeMismatch,
+ ResponseHashMismatch,
InvalidDescriptor,
DescriptorMismatch,
RetrievedBytesMismatch,
@@ -976,7 +1060,10 @@ impl BlossomError {
BlossomErrorKind::HttpStatus => "blossom_http_status",
BlossomErrorKind::UnsafeRedirect => "blossom_unsafe_redirect",
BlossomErrorKind::RedirectLimit => "blossom_redirect_limit",
+ BlossomErrorKind::ContentEncodingDenied => "blossom_content_encoding_denied",
BlossomErrorKind::ResponseTooLarge => "blossom_response_too_large",
+ BlossomErrorKind::ResponseSizeMismatch => "blossom_response_size_mismatch",
+ BlossomErrorKind::ResponseHashMismatch => "blossom_response_hash_mismatch",
BlossomErrorKind::InvalidDescriptor => "blossom_invalid_descriptor",
BlossomErrorKind::DescriptorMismatch => "blossom_descriptor_mismatch",
BlossomErrorKind::RetrievedBytesMismatch => "blossom_retrieved_bytes_mismatch",
@@ -1036,6 +1123,93 @@ pub struct BlossomUploadReceipt {
verified_at_unix_ms: u64,
}
+/// Exact verified image bytes returned by a bounded BUD-01 retrieval.
+#[cfg(feature = "blossom")]
+#[derive(Clone)]
+pub struct BlossomInboundReceipt {
+ final_url: BlobUrl,
+ commitment: radroots_blossom::descriptor::ByteCommitment,
+ dimensions: BlossomImageDimensions,
+ bytes: Arc<[u8]>,
+ config_fingerprint: BlossomConfigFingerprint,
+ attempts: u8,
+ verified_at_unix_ms: u64,
+}
+
+#[cfg(feature = "blossom")]
+impl std::fmt::Debug for BlossomInboundReceipt {
+ fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
+ formatter
+ .debug_struct("BlossomInboundReceipt")
+ .field("final_url", &self.final_url)
+ .field("commitment", &self.commitment)
+ .field("dimensions", &self.dimensions)
+ .field("bytes", &"<redacted>")
+ .field("config_fingerprint", &self.config_fingerprint)
+ .field("attempts", &self.attempts)
+ .field("verified_at_unix_ms", &self.verified_at_unix_ms)
+ .finish()
+ }
+}
+
+#[cfg(feature = "blossom")]
+impl BlossomInboundReceipt {
+ pub(crate) fn new(
+ final_url: BlobUrl,
+ commitment: radroots_blossom::descriptor::ByteCommitment,
+ dimensions: BlossomImageDimensions,
+ bytes: Arc<[u8]>,
+ config_fingerprint: BlossomConfigFingerprint,
+ attempts: u8,
+ verified_at_unix_ms: u64,
+ ) -> Self {
+ Self {
+ final_url,
+ commitment,
+ dimensions,
+ bytes,
+ config_fingerprint,
+ attempts,
+ verified_at_unix_ms,
+ }
+ }
+
+ #[must_use]
+ pub const fn final_url(&self) -> &BlobUrl {
+ &self.final_url
+ }
+
+ #[must_use]
+ pub const fn commitment(&self) -> &radroots_blossom::descriptor::ByteCommitment {
+ &self.commitment
+ }
+
+ #[must_use]
+ pub const fn dimensions(&self) -> BlossomImageDimensions {
+ self.dimensions
+ }
+
+ #[must_use]
+ pub fn bytes(&self) -> &[u8] {
+ self.bytes.as_ref()
+ }
+
+ #[must_use]
+ pub const fn config_fingerprint(&self) -> BlossomConfigFingerprint {
+ self.config_fingerprint
+ }
+
+ #[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
+ }
+}
+
#[cfg(feature = "blossom")]
impl BlossomUploadReceipt {
pub(crate) const fn new(
@@ -1279,6 +1453,39 @@ impl BlossomSlot {
}
}
+ /// Retrieves and verifies one immutable BUD-01 image under the exact
+ /// configured DNS, TLS, redirect, retry, and byte limits.
+ pub async fn retrieve(
+ &self,
+ request: BlossomInboundRequest,
+ cancellation: BlossomCancellation,
+ ) -> Result<BlossomInboundReceipt, BlossomError> {
+ let config = self
+ .snapshot()
+ .ok_or_else(|| BlossomError::configuration(BlossomErrorKind::EndpointNotConfigured))?;
+ let fingerprint = config.fingerprint();
+ if config.profile().endpoint_for_blob(request.url()).is_none() {
+ return Err(BlossomError::configuration(
+ BlossomErrorKind::EndpointNotConfigured,
+ ));
+ }
+ let result = crate::adapters::blossom::retrieve(config, request, cancellation).await;
+ match result {
+ Ok(receipt) => {
+ self.record_evidence(fingerprint, BlossomPhase::Verification, false, |evidence| {
+ evidence.record_success(BlossomEvidenceState::RetrievalVerified, None);
+ })?;
+ Ok(receipt)
+ }
+ Err(error) => {
+ self.record_evidence(fingerprint, BlossomPhase::Verification, false, |evidence| {
+ evidence.record_failure(&error);
+ })?;
+ Err(error)
+ }
+ }
+ }
+
fn validate_transaction(
&self,
transaction: &BlossomUploadTransaction,
@@ -1414,7 +1621,9 @@ fn blossom_authority_accepts_address(authority: BlossomEndpointAuthority, addres
}
#[cfg(feature = "blossom")]
-fn canonical_image_extension(media_type: &MediaType) -> Result<&'static str, BlossomError> {
+pub(crate) fn canonical_image_extension(
+ media_type: &MediaType,
+) -> Result<&'static str, BlossomError> {
match media_type.as_str() {
"image/png" => Ok("png"),
"image/jpeg" => Ok("jpg"),
@@ -2571,6 +2780,8 @@ mod tests {
assert!(BlossomImageDimensions::new(0, 1).is_err());
assert!(BlossomImageDimensions::new(1, 0).is_err());
+ assert!(BlossomImageDimensions::new(16_385, 1).is_err());
+ assert!(BlossomImageDimensions::new(10_001, 10_000).is_err());
let dimensions = BlossomImageDimensions::new(2, 3).unwrap();
assert_eq!(dimensions.width(), 2);
assert_eq!(dimensions.height(), 3);
@@ -2665,10 +2876,22 @@ mod tests {
(BlossomErrorKind::UnsafeRedirect, "blossom_unsafe_redirect"),
(BlossomErrorKind::RedirectLimit, "blossom_redirect_limit"),
(
+ BlossomErrorKind::ContentEncodingDenied,
+ "blossom_content_encoding_denied",
+ ),
+ (
BlossomErrorKind::ResponseTooLarge,
"blossom_response_too_large",
),
(
+ BlossomErrorKind::ResponseSizeMismatch,
+ "blossom_response_size_mismatch",
+ ),
+ (
+ BlossomErrorKind::ResponseHashMismatch,
+ "blossom_response_hash_mismatch",
+ ),
+ (
BlossomErrorKind::InvalidDescriptor,
"blossom_invalid_descriptor",
),
diff --git a/crates/sdk/tests/public_api.rs b/crates/sdk/tests/public_api.rs
@@ -85,6 +85,8 @@ fn public_native_type_snapshot_uses_contextual_names() {
type_name::<radroots_sdk::transport::BlossomEvidenceState>(),
type_name::<radroots_sdk::transport::BlossomHostKind>(),
type_name::<radroots_sdk::transport::BlossomImageDimensions>(),
+ type_name::<radroots_sdk::transport::BlossomInboundReceipt>(),
+ type_name::<radroots_sdk::transport::BlossomInboundRequest>(),
type_name::<radroots_sdk::transport::BlossomPhase>(),
type_name::<radroots_sdk::transport::BlossomProfile>(),
type_name::<radroots_sdk::transport::BlossomSlot>(),
@@ -171,7 +173,7 @@ const fn expected_public_type_count() -> usize {
#[cfg(feature = "nostr")]
let count = count + 1;
#[cfg(feature = "blossom")]
- let count = count + 20;
+ let count = count + 22;
#[cfg(feature = "radrootsd")]
let count = count + 5;
count