commit a7544771421c9aae94da3647afd2382b60e7d29a
parent 28d6f6b901972bd524554b51531f23745c1d2b4f
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:
7 files changed, 882 insertions(+), 9 deletions(-)
diff --git a/core/crates/tera_core/Cargo.toml b/core/crates/tera_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/core/crates/tera_core/src/runtime/builder.rs b/core/crates/tera_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/core/crates/tera_core/src/runtime/mod.rs b/core/crates/tera_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/core/crates/tera_core/src/runtime/product_surface.rs b/core/crates/tera_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/core/crates/tera_core/src/runtime/product_surface/media.rs b/core/crates/tera_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/core/crates/tera_core/src/runtime/product_surface/outbox.rs b/core/crates/tera_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/core/crates/tera_core/src/runtime/product_surface/today.rs b/core/crates/tera_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");