lib

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

commit f7218237bb70e97fc9cd8f902dc8cffa9706d667
parent 47d1a4f28ae5242e96d6c8ce9db6629c6d136a22
Author: triesap <tyson@radroots.org>
Date:   Mon, 10 Aug 2026 17:32:12 +0000

feat: persist verified inbound media receipts

Separate signed structural references from retrieval state and exact-byte receipts.

Persist a configuration-scoped, content-addressed, bounded LRU cache and fail closed across eviction, corruption, metadata drift, and legacy migration.

Diffstat:
MCargo.lock | 1+
Mcrates/mobile_core/Cargo.toml | 1+
Mcrates/mobile_core/src/runtime/product_surface.rs | 12+++++++++---
Acrates/mobile_core/src/runtime/product_surface/media.rs | 1241+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcrates/mobile_core/src/runtime/product_surface/model.rs | 38+-------------------------------------
Mcrates/mobile_core/src/runtime/product_surface/projection.rs | 66+++++++++++++++++++++++++++++++++++-------------------------------
Mcrates/mobile_core/src/runtime/product_surface/today.rs | 1061+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++----------
Mcrates/mobile_ffi/src/dto.rs | 35++++++++++++++++++-----------------
Mcrates/mobile_ffi/src/error.rs | 1+
9 files changed, 2245 insertions(+), 211 deletions(-)

diff --git a/Cargo.lock b/Cargo.lock @@ -3599,6 +3599,7 @@ dependencies = [ "tempfile", "thiserror 1.0.69", "tokio", + "url", "uuid", ] diff --git a/crates/mobile_core/Cargo.toml b/crates/mobile_core/Cargo.toml @@ -50,6 +50,7 @@ sha2 = { workspace = true } thiserror = { workspace = true } tokio = { workspace = true, optional = true, features = ["sync"] } uuid = { workspace = true, features = ["v4"] } +url = { workspace = true } [dev-dependencies] nostr = { workspace = true, features = ["std"] } diff --git a/crates/mobile_core/src/runtime/product_surface.rs b/crates/mobile_core/src/runtime/product_surface.rs @@ -8,6 +8,7 @@ mod authoring; mod context; mod cursor; mod identity; +mod media; mod model; #[cfg(feature = "mobile-social")] mod outbox; @@ -27,12 +28,17 @@ pub use context::{ }; pub use cursor::{CursorError, CursorScope, TodayCursor, TodayCursorPosition}; pub use identity::{CARD_ID_SCHEMA_VERSION, CardId, CardIdError, CardSourceIdentity}; +pub use media::{ + MediaReference, Phase1InboundMediaError, Phase1InboundMediaFailure, Phase1InboundMediaPending, + Phase1InboundMediaState, Phase1MediaArtifactId, Phase1MediaCacheIndex, Phase1MediaCachePolicy, + Phase1MediaCacheStatus, Phase1MediaConfigurationFingerprint, Phase1StructuralMediaReference, + Phase1VerifiedMediaReceipt, +}; pub use model::{ AddCommandType, CANONICAL_ADD_COMMAND_TYPES, CANONICAL_CARD_ADD_PARITY, CANONICAL_TODAY_CARD_TYPES, CardAddParity, CardLifecycleState, ClassifiedCard, - LocalAuthorOverlay, MeSnapshot, MediaReference, MediaVerificationState, ProfileSummary, - SearchResult, SearchResultType, SupportingProfile, ThreadEntry, ThreadReference, TodayCard, - TodayCardType, TodayPage, + LocalAuthorOverlay, MeSnapshot, ProfileSummary, SearchResult, SearchResultType, + SupportingProfile, ThreadEntry, ThreadReference, TodayCard, TodayCardType, TodayPage, }; #[cfg(feature = "mobile-social")] pub use outbox::{ diff --git a/crates/mobile_core/src/runtime/product_surface/media.rs b/crates/mobile_core/src/runtime/product_surface/media.rs @@ -0,0 +1,1241 @@ +//! Typed inbound-media trust, receipt, and bounded cache metadata. +//! +//! Structural Nostr references are intentionally distinct from locally +//! verified artifacts. A caller cannot represent renderable media with a URL +//! and a boolean: the `Verified` state always contains a receipt derived from +//! an actual byte commitment and bound to the active retrieval configuration. + +use std::collections::BTreeMap; + +use radroots_blossom::{BlobUrl, MediaType, Sha256, descriptor::ByteCommitment}; +use serde::{Deserialize, Serialize}; +use sha2::{Digest, Sha256 as Sha256Hasher}; +use thiserror::Error; + +const MEDIA_REFERENCE_SCHEMA_VERSION: u16 = 1; +const MEDIA_RECEIPT_SCHEMA_VERSION: u16 = 1; +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_REFERENCE_FINGERPRINT_DOMAIN: &[u8] = b"radroots.inbound-media-reference.v1\0"; + +#[derive(Clone, Copy, Debug, Deserialize, Eq, Hash, Ord, PartialEq, PartialOrd, Serialize)] +#[serde(transparent)] +pub struct Phase1MediaConfigurationFingerprint([u8; 32]); + +impl Phase1MediaConfigurationFingerprint { + pub fn new(value: [u8; 32]) -> Result<Self, Phase1InboundMediaError> { + (value != [0; 32]) + .then_some(Self(value)) + .ok_or(Phase1InboundMediaError::InvalidConfiguration) + } + + pub fn parse(value: &str) -> Result<Self, Phase1InboundMediaError> { + let decoded = + hex::decode(value).map_err(|_| Phase1InboundMediaError::InvalidConfiguration)?; + let bytes: [u8; 32] = decoded + .try_into() + .map_err(|_| Phase1InboundMediaError::InvalidConfiguration)?; + Self::new(bytes) + } + + pub const fn as_bytes(&self) -> &[u8; 32] { + &self.0 + } + + pub fn to_hex(self) -> String { + hex::encode(self.0) + } + + fn validate(self) -> Result<(), Phase1InboundMediaError> { + Self::new(self.0).map(|_| ()) + } +} + +#[derive(Clone, Copy, Debug, Deserialize, Eq, Hash, Ord, PartialEq, PartialOrd, Serialize)] +#[serde(transparent)] +pub struct Phase1MediaArtifactId([u8; 32]); + +impl Phase1MediaArtifactId { + pub fn parse(value: &str) -> Result<Self, Phase1InboundMediaError> { + let hash = Sha256::from_hex(value).map_err(|_| Phase1InboundMediaError::InvalidDigest)?; + Ok(Self(*hash.as_bytes())) + } + + pub const fn from_sha256(value: Sha256) -> Self { + Self(*value.as_bytes()) + } + + pub const fn as_bytes(&self) -> &[u8; 32] { + &self.0 + } + + pub fn to_hex(self) -> String { + hex::encode(self.0) + } +} + +/// Signed-event media facts. These facts do not imply that any bytes were +/// fetched, trusted, stored, or rendered. +#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)] +#[serde(deny_unknown_fields, rename_all = "camelCase")] +pub struct Phase1StructuralMediaReference { + schema_version: u16, + source_url: String, + expected_sha256: Option<String>, + expected_media_type: Option<String>, + expected_width: Option<u32>, + expected_height: Option<u32>, + expected_byte_size: Option<u64>, + alt: Option<String>, + fingerprint: [u8; 32], +} + +impl Phase1StructuralMediaReference { + #[allow(clippy::too_many_arguments)] + pub fn new( + source_url: impl Into<String>, + expected_sha256: Option<String>, + expected_media_type: Option<String>, + expected_width: Option<u32>, + expected_height: Option<u32>, + expected_byte_size: Option<u64>, + alt: Option<String>, + ) -> Result<Self, Phase1InboundMediaError> { + let source_url = source_url.into(); + validate_url_text(&source_url)?; + let parsed = + url::Url::parse(&source_url).map_err(|_| Phase1InboundMediaError::InvalidReference)?; + if !matches!(parsed.scheme(), "http" | "https") + || parsed.host_str().is_none() + || !parsed.username().is_empty() + || parsed.password().is_some() + || parsed.as_str() != source_url + { + return Err(Phase1InboundMediaError::InvalidReference); + } + let path_digest = BlobUrl::parse(&source_url) + .ok() + .map(|value| value.hash_path().hash().to_hex()); + let expected_sha256 = match expected_sha256 { + Some(value) => { + let digest = + Sha256::from_hex(&value).map_err(|_| Phase1InboundMediaError::InvalidDigest)?; + if digest.to_hex() != value + || path_digest + .as_deref() + .is_some_and(|path| digest.to_hex() != path) + { + return Err(Phase1InboundMediaError::MetadataMismatch); + } + Some(value) + } + None => path_digest, + }; + let expected_media_type = expected_media_type + .map(|value| { + MediaType::parse(&value) + .map(|parsed| parsed.to_string()) + .map_err(|_| Phase1InboundMediaError::InvalidMediaType) + }) + .transpose()?; + validate_dimensions(expected_width, expected_height)?; + if expected_byte_size == Some(0) { + return Err(Phase1InboundMediaError::InvalidByteSize); + } + if alt.as_deref().is_some_and(|value| { + value.len() > MEDIA_ALT_MAX_BYTES || value.chars().any(char::is_control) + }) { + return Err(Phase1InboundMediaError::InvalidAlt); + } + let mut value = Self { + schema_version: MEDIA_REFERENCE_SCHEMA_VERSION, + source_url: parsed.to_string(), + expected_sha256, + expected_media_type, + expected_width, + expected_height, + expected_byte_size, + alt, + fingerprint: [0; 32], + }; + value.fingerprint = value.derive_fingerprint(); + Ok(value) + } + + fn validate(&self) -> Result<(), Phase1InboundMediaError> { + if self.schema_version != MEDIA_REFERENCE_SCHEMA_VERSION { + return Err(Phase1InboundMediaError::UnsupportedSchema); + } + let canonical = Self::new( + self.source_url.clone(), + self.expected_sha256.clone(), + self.expected_media_type.clone(), + self.expected_width, + self.expected_height, + self.expected_byte_size, + self.alt.clone(), + )?; + (canonical == *self) + .then_some(()) + .ok_or(Phase1InboundMediaError::CorruptState) + } + + fn derive_fingerprint(&self) -> [u8; 32] { + let mut digest = Sha256Hasher::new(); + digest.update(MEDIA_REFERENCE_FINGERPRINT_DOMAIN); + digest.update(self.source_url.as_bytes()); + update_optional(&mut digest, self.expected_sha256.as_deref()); + update_optional(&mut digest, self.expected_media_type.as_deref()); + update_optional_u64(&mut digest, self.expected_width.map(u64::from)); + update_optional_u64(&mut digest, self.expected_height.map(u64::from)); + update_optional_u64(&mut digest, self.expected_byte_size); + update_optional(&mut digest, self.alt.as_deref()); + digest.finalize().into() + } + + pub fn source_url(&self) -> &str { + self.source_url.as_str() + } + + pub fn expected_sha256(&self) -> Option<&str> { + self.expected_sha256.as_deref() + } + + pub fn expected_media_type(&self) -> Option<&str> { + self.expected_media_type.as_deref() + } + + pub const fn expected_width(&self) -> Option<u32> { + self.expected_width + } + + pub const fn expected_height(&self) -> Option<u32> { + self.expected_height + } + + pub const fn expected_byte_size(&self) -> Option<u64> { + self.expected_byte_size + } + + pub fn alt(&self) -> Option<&str> { + self.alt.as_deref() + } + + pub const fn fingerprint(&self) -> &[u8; 32] { + &self.fingerprint + } +} + +#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)] +#[serde(deny_unknown_fields, rename_all = "camelCase")] +pub struct Phase1InboundMediaPending { + operation_id: [u8; 16], + configuration: Phase1MediaConfigurationFingerprint, + started_at_unix_ms: u64, +} + +impl Phase1InboundMediaPending { + pub fn new( + operation_id: [u8; 16], + configuration: Phase1MediaConfigurationFingerprint, + started_at_unix_ms: u64, + ) -> Result<Self, Phase1InboundMediaError> { + configuration.validate()?; + if operation_id == [0; 16] || started_at_unix_ms == 0 { + return Err(Phase1InboundMediaError::InvalidOperation); + } + Ok(Self { + operation_id, + configuration, + started_at_unix_ms, + }) + } + + pub const fn operation_id(&self) -> &[u8; 16] { + &self.operation_id + } + + pub const fn configuration(&self) -> Phase1MediaConfigurationFingerprint { + self.configuration + } + + pub const fn started_at_unix_ms(&self) -> u64 { + self.started_at_unix_ms + } + + fn validate(&self) -> Result<(), Phase1InboundMediaError> { + Self::new( + self.operation_id, + self.configuration, + self.started_at_unix_ms, + ) + .map(|_| ()) + } +} + +#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)] +#[serde(deny_unknown_fields, rename_all = "camelCase")] +pub struct Phase1InboundMediaFailure { + operation_id: [u8; 16], + safe_code: String, + retryable: bool, + failed_at_unix_ms: u64, +} + +impl Phase1InboundMediaFailure { + pub fn new( + operation_id: [u8; 16], + safe_code: impl Into<String>, + retryable: bool, + failed_at_unix_ms: u64, + ) -> Result<Self, Phase1InboundMediaError> { + let safe_code = safe_code.into(); + if operation_id == [0; 16] + || failed_at_unix_ms == 0 + || safe_code.is_empty() + || safe_code.len() > MEDIA_FAILURE_CODE_MAX_BYTES + || !safe_code + .bytes() + .all(|byte| byte.is_ascii_lowercase() || byte.is_ascii_digit() || byte == b'_') + { + return Err(Phase1InboundMediaError::InvalidFailure); + } + Ok(Self { + operation_id, + safe_code, + retryable, + failed_at_unix_ms, + }) + } + + pub fn safe_code(&self) -> &str { + self.safe_code.as_str() + } + + pub const fn retryable(&self) -> bool { + self.retryable + } + + fn validate(&self) -> Result<(), Phase1InboundMediaError> { + Self::new( + self.operation_id, + self.safe_code.clone(), + self.retryable, + self.failed_at_unix_ms, + ) + .map(|_| ()) + } +} + +/// Exact-byte verification evidence. Construction requires a byte commitment, +/// binds every signed expected field, and derives the artifact identity from +/// the observed digest rather than caller input. +#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)] +#[serde(deny_unknown_fields, rename_all = "camelCase")] +pub struct Phase1VerifiedMediaReceipt { + schema_version: u16, + reference_fingerprint: [u8; 32], + source_url: String, + canonical_final_url: String, + expected_sha256: String, + observed_sha256: String, + byte_size: u64, + media_type: String, + extension: String, + width: u32, + height: u32, + artifact_id: Phase1MediaArtifactId, + configuration: Phase1MediaConfigurationFingerprint, + verified_at_unix_ms: u64, +} + +impl Phase1VerifiedMediaReceipt { + pub fn from_commitment( + reference: &Phase1StructuralMediaReference, + canonical_final_url: BlobUrl, + commitment: &ByteCommitment, + width: u32, + height: u32, + configuration: Phase1MediaConfigurationFingerprint, + verified_at_unix_ms: u64, + ) -> Result<Self, Phase1InboundMediaError> { + reference.validate()?; + configuration.validate()?; + if verified_at_unix_ms == 0 { + return Err(Phase1InboundMediaError::InvalidVerificationTime); + } + BlobUrl::parse(reference.source_url()) + .and_then(BlobUrl::approve) + .map_err(|_| Phase1InboundMediaError::InvalidReference)?; + canonical_final_url + .clone() + .approve() + .map_err(|_| Phase1InboundMediaError::InvalidReference)?; + validate_dimensions(Some(width), Some(height))?; + let observed_sha256 = commitment.sha256().to_hex(); + let expected_sha256 = reference + .expected_sha256() + .ok_or(Phase1InboundMediaError::MissingDigest)?; + if observed_sha256 != expected_sha256 + || reference + .expected_byte_size() + .is_some_and(|value| value != commitment.size()) + || reference + .expected_media_type() + .is_some_and(|value| value != commitment.media_type().to_string()) + || reference + .expected_width() + .is_some_and(|value| value != width) + || reference + .expected_height() + .is_some_and(|value| value != height) + { + return Err(Phase1InboundMediaError::MetadataMismatch); + } + let extension = canonical_final_url + .hash_path() + .extension() + .ok_or(Phase1InboundMediaError::InvalidReference)? + .as_str() + .to_owned(); + if canonical_final_url.hash_path().hash() != commitment.sha256() { + return Err(Phase1InboundMediaError::MetadataMismatch); + } + let receipt = Self { + schema_version: MEDIA_RECEIPT_SCHEMA_VERSION, + reference_fingerprint: *reference.fingerprint(), + source_url: reference.source_url().to_owned(), + canonical_final_url: canonical_final_url.to_string(), + expected_sha256: expected_sha256.to_owned(), + observed_sha256, + byte_size: commitment.size(), + media_type: commitment.media_type().to_string(), + extension, + width, + height, + artifact_id: Phase1MediaArtifactId::from_sha256(commitment.sha256()), + configuration, + verified_at_unix_ms, + }; + receipt.validate(reference)?; + Ok(receipt) + } + + fn validate( + &self, + reference: &Phase1StructuralMediaReference, + ) -> Result<(), Phase1InboundMediaError> { + self.validate_intrinsic()?; + if self.reference_fingerprint != *reference.fingerprint() + || self.source_url != reference.source_url() + || reference.expected_sha256() != Some(self.expected_sha256.as_str()) + { + return Err(Phase1InboundMediaError::CorruptReceipt); + } + BlobUrl::parse(reference.source_url()) + .and_then(BlobUrl::approve) + .map_err(|_| Phase1InboundMediaError::CorruptReceipt)?; + if reference + .expected_byte_size() + .is_some_and(|value| value != self.byte_size) + || reference + .expected_media_type() + .is_some_and(|value| value != self.media_type) + || reference + .expected_width() + .is_some_and(|value| value != self.width) + || reference + .expected_height() + .is_some_and(|value| value != self.height) + { + return Err(Phase1InboundMediaError::CorruptReceipt); + } + Ok(()) + } + + fn validate_intrinsic(&self) -> Result<(), Phase1InboundMediaError> { + if self.schema_version != MEDIA_RECEIPT_SCHEMA_VERSION + || self.expected_sha256 != self.observed_sha256 + || self.expected_sha256 != self.artifact_id.to_hex() + || self.byte_size == 0 + || self.verified_at_unix_ms == 0 + { + return Err(Phase1InboundMediaError::CorruptReceipt); + } + self.configuration + .validate() + .map_err(|_| Phase1InboundMediaError::CorruptReceipt)?; + let final_url = BlobUrl::parse(&self.canonical_final_url) + .map_err(|_| Phase1InboundMediaError::CorruptReceipt)?; + final_url + .clone() + .approve() + .map_err(|_| Phase1InboundMediaError::CorruptReceipt)?; + if final_url.to_string() != self.canonical_final_url + || final_url.hash_path().hash().to_hex() != self.observed_sha256 + || final_url + .hash_path() + .extension() + .is_none_or(|value| value.as_str() != self.extension) + || !canonical_media_type(&self.media_type) + || validate_dimensions(Some(self.width), Some(self.height)).is_err() + { + return Err(Phase1InboundMediaError::CorruptReceipt); + } + Ok(()) + } + + pub const fn artifact_id(&self) -> Phase1MediaArtifactId { + self.artifact_id + } + + pub const fn configuration(&self) -> Phase1MediaConfigurationFingerprint { + self.configuration + } + + pub fn canonical_final_url(&self) -> &str { + self.canonical_final_url.as_str() + } + + pub fn observed_sha256(&self) -> &str { + self.observed_sha256.as_str() + } + + pub const fn byte_size(&self) -> u64 { + self.byte_size + } + + pub fn media_type(&self) -> &str { + self.media_type.as_str() + } + + pub fn extension(&self) -> &str { + self.extension.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)] +#[serde(rename_all = "camelCase", tag = "state", content = "evidence")] +pub enum Phase1InboundMediaState { + #[default] + Unavailable, + Pending(Phase1InboundMediaPending), + Failed(Phase1InboundMediaFailure), + Verified(Box<Phase1VerifiedMediaReceipt>), +} + +/// Public media model: signed structure plus local retrieval evidence. +#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)] +#[serde(deny_unknown_fields, rename_all = "camelCase")] +pub struct MediaReference { + structural: Phase1StructuralMediaReference, + retrieval: Phase1InboundMediaState, +} + +impl MediaReference { + pub fn new( + structural: Phase1StructuralMediaReference, + ) -> Result<Self, Phase1InboundMediaError> { + structural.validate()?; + Ok(Self { + structural, + retrieval: Phase1InboundMediaState::Unavailable, + }) + } + + pub(crate) fn legacy_unavailable( + source_url: String, + expected_sha256: Option<String>, + expected_media_type: Option<String>, + expected_width: Option<u32>, + expected_height: Option<u32>, + expected_byte_size: Option<u64>, + alt: Option<String>, + ) -> Result<Self, Phase1InboundMediaError> { + Self::new(Phase1StructuralMediaReference::new( + source_url, + expected_sha256, + expected_media_type, + expected_width, + expected_height, + expected_byte_size, + alt, + )?) + } + + pub fn structural(&self) -> &Phase1StructuralMediaReference { + &self.structural + } + + pub const fn retrieval(&self) -> &Phase1InboundMediaState { + &self.retrieval + } + + pub(crate) fn validate(&self) -> Result<(), Phase1InboundMediaError> { + self.structural.validate()?; + match &self.retrieval { + Phase1InboundMediaState::Unavailable => Ok(()), + Phase1InboundMediaState::Pending(value) => value.validate(), + Phase1InboundMediaState::Failed(value) => value.validate(), + Phase1InboundMediaState::Verified(value) => value.validate(&self.structural), + } + } + + pub(crate) fn restore( + &mut self, + retrieval: Phase1InboundMediaState, + cache: &Phase1MediaCacheIndex, + ) -> Result<(), Phase1InboundMediaError> { + let mut candidate = self.clone(); + candidate.retrieval = retrieval; + candidate.validate()?; + if let Phase1InboundMediaState::Verified(receipt) = &candidate.retrieval + && !cache.contains(receipt) + { + candidate.retrieval = Phase1InboundMediaState::Unavailable; + } + *self = candidate; + Ok(()) + } + + pub fn begin( + &mut self, + pending: Phase1InboundMediaPending, + ) -> Result<(), Phase1InboundMediaError> { + self.structural.validate()?; + pending.validate()?; + self.retrieval = Phase1InboundMediaState::Pending(pending); + Ok(()) + } + + pub fn fail( + &mut self, + failure: Phase1InboundMediaFailure, + ) -> Result<(), Phase1InboundMediaError> { + failure.validate()?; + match &self.retrieval { + Phase1InboundMediaState::Pending(pending) + if pending.operation_id == failure.operation_id => + { + self.retrieval = Phase1InboundMediaState::Failed(failure); + Ok(()) + } + _ => Err(Phase1InboundMediaError::OperationMismatch), + } + } + + pub fn verify( + &mut self, + operation_id: [u8; 16], + receipt: Phase1VerifiedMediaReceipt, + ) -> Result<(), Phase1InboundMediaError> { + let Phase1InboundMediaState::Pending(pending) = &self.retrieval else { + return Err(Phase1InboundMediaError::OperationMismatch); + }; + if pending.operation_id != operation_id || pending.configuration != receipt.configuration { + return Err(Phase1InboundMediaError::OperationMismatch); + } + receipt.validate(&self.structural)?; + self.retrieval = Phase1InboundMediaState::Verified(Box::new(receipt)); + Ok(()) + } + + pub fn invalidate(&mut self) -> Option<Phase1MediaArtifactId> { + let artifact = match &self.retrieval { + Phase1InboundMediaState::Verified(receipt) => Some(receipt.artifact_id), + _ => None, + }; + self.retrieval = Phase1InboundMediaState::Unavailable; + artifact + } + + pub fn is_renderable_with( + &self, + cache: &Phase1MediaCacheIndex, + configuration: Phase1MediaConfigurationFingerprint, + ) -> bool { + match &self.retrieval { + Phase1InboundMediaState::Verified(receipt) + if receipt.configuration == configuration => + { + cache.contains(receipt) + } + _ => false, + } + } +} + +#[derive(Clone, Copy, Debug, Deserialize, Eq, PartialEq, Serialize)] +#[serde(deny_unknown_fields, rename_all = "camelCase")] +pub struct Phase1MediaCachePolicy { + max_bytes: u64, + max_artifacts: u32, +} + +impl Phase1MediaCachePolicy { + pub fn new(max_bytes: u64, max_artifacts: u32) -> Result<Self, Phase1InboundMediaError> { + if max_bytes == 0 || max_artifacts == 0 { + return Err(Phase1InboundMediaError::InvalidCachePolicy); + } + Ok(Self { + max_bytes, + max_artifacts, + }) + } + + pub const fn max_bytes(&self) -> u64 { + self.max_bytes + } + + pub const fn max_artifacts(&self) -> u32 { + self.max_artifacts + } +} + +impl Default for Phase1MediaCachePolicy { + fn default() -> Self { + Self { + max_bytes: 256 * 1024 * 1024, + max_artifacts: 2_000, + } + } +} + +#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)] +#[serde(deny_unknown_fields, rename_all = "camelCase")] +struct Phase1MediaCacheEntry { + artifact_id: Phase1MediaArtifactId, + byte_size: u64, + media_type: String, + extension: String, + width: u32, + height: u32, + cached_at_unix_ms: u64, + last_accessed_at_unix_ms: u64, +} + +impl Phase1MediaCacheEntry { + fn from_receipt( + receipt: &Phase1VerifiedMediaReceipt, + cached_at_unix_ms: u64, + ) -> Result<Self, Phase1InboundMediaError> { + if cached_at_unix_ms < receipt.verified_at_unix_ms { + return Err(Phase1InboundMediaError::InvalidCacheObservation); + } + Ok(Self { + artifact_id: receipt.artifact_id, + byte_size: receipt.byte_size, + media_type: receipt.media_type.clone(), + extension: receipt.extension.clone(), + width: receipt.width, + height: receipt.height, + cached_at_unix_ms, + last_accessed_at_unix_ms: cached_at_unix_ms, + }) + } + + fn matches(&self, receipt: &Phase1VerifiedMediaReceipt) -> bool { + self.artifact_id == receipt.artifact_id + && self.byte_size == receipt.byte_size + && self.media_type == receipt.media_type + && self.extension == receipt.extension + && self.width == receipt.width + && self.height == receipt.height + } +} + +#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)] +#[serde(deny_unknown_fields, rename_all = "camelCase")] +pub struct Phase1MediaCacheIndex { + schema_version: u16, + configuration: Option<Phase1MediaConfigurationFingerprint>, + entries: BTreeMap<String, Phase1MediaCacheEntry>, +} + +impl Default for Phase1MediaCacheIndex { + fn default() -> Self { + Self { + schema_version: MEDIA_CACHE_SCHEMA_VERSION, + configuration: None, + entries: BTreeMap::new(), + } + } +} + +impl Phase1MediaCacheIndex { + pub fn admit( + &mut self, + receipt: &Phase1VerifiedMediaReceipt, + policy: Phase1MediaCachePolicy, + cached_at_unix_ms: u64, + ) -> Result<Vec<Phase1MediaArtifactId>, Phase1InboundMediaError> { + self.validate()?; + receipt.validate_intrinsic()?; + if receipt.byte_size > policy.max_bytes { + return Err(Phase1InboundMediaError::CacheQuotaExceeded); + } + if self + .configuration + .is_some_and(|value| value != receipt.configuration) + { + return Err(Phase1InboundMediaError::ConfigurationMismatch); + } + self.configuration = Some(receipt.configuration); + let key = receipt.artifact_id.to_hex(); + let entry = Phase1MediaCacheEntry::from_receipt(receipt, cached_at_unix_ms)?; + if self + .entries + .get(&key) + .is_some_and(|existing| !existing.matches(receipt)) + { + return Err(Phase1InboundMediaError::ArtifactCollision); + } + self.entries.insert(key, entry); + let mut evicted = Vec::new(); + while self.entries.len() > policy.max_artifacts as usize + || self.total_bytes()? > policy.max_bytes + { + let oldest = self + .entries + .iter() + .min_by_key(|(key, entry)| (entry.last_accessed_at_unix_ms, key.as_str())) + .map(|(key, _)| key.clone()) + .ok_or(Phase1InboundMediaError::CorruptState)?; + let removed = self + .entries + .remove(&oldest) + .ok_or(Phase1InboundMediaError::CorruptState)?; + evicted.push(removed.artifact_id); + } + Ok(evicted) + } + + pub fn contains(&self, receipt: &Phase1VerifiedMediaReceipt) -> bool { + self.schema_version == MEDIA_CACHE_SCHEMA_VERSION + && self.configuration == Some(receipt.configuration) + && self + .entries + .get(&receipt.artifact_id.to_hex()) + .is_some_and(|entry| entry.matches(receipt)) + } + + pub fn touch( + &mut self, + artifact_id: Phase1MediaArtifactId, + observed_at_unix_ms: u64, + ) -> Result<bool, Phase1InboundMediaError> { + if observed_at_unix_ms == 0 { + return Err(Phase1InboundMediaError::InvalidCacheObservation); + } + let Some(entry) = self.entries.get_mut(&artifact_id.to_hex()) else { + return Ok(false); + }; + entry.last_accessed_at_unix_ms = entry.last_accessed_at_unix_ms.max(observed_at_unix_ms); + Ok(true) + } + + pub fn invalidate_artifact(&mut self, artifact_id: Phase1MediaArtifactId) -> bool { + self.entries.remove(&artifact_id.to_hex()).is_some() + } + + pub fn invalidate_configuration( + &mut self, + configuration: Phase1MediaConfigurationFingerprint, + ) -> Vec<Phase1MediaArtifactId> { + if self + .configuration + .is_none_or(|current| current == configuration) + { + self.configuration = Some(configuration); + return Vec::new(); + } + let removed = self + .entries + .values() + .map(|entry| entry.artifact_id) + .collect(); + self.entries.clear(); + self.configuration = Some(configuration); + removed + } + + pub fn artifact_count(&self) -> u32 { + self.entries.len().try_into().unwrap_or(u32::MAX) + } + + pub fn total_bytes(&self) -> Result<u64, Phase1InboundMediaError> { + self.entries.values().try_fold(0_u64, |total, entry| { + total + .checked_add(entry.byte_size) + .ok_or(Phase1InboundMediaError::CorruptState) + }) + } + + fn validate(&self) -> Result<(), Phase1InboundMediaError> { + if self.schema_version != MEDIA_CACHE_SCHEMA_VERSION + || (self.configuration.is_none() && !self.entries.is_empty()) + || self.entries.iter().any(|(key, entry)| { + key != &entry.artifact_id.to_hex() + || key != &hex::encode(entry.artifact_id.as_bytes()) + || entry.byte_size == 0 + || entry.cached_at_unix_ms == 0 + || entry.last_accessed_at_unix_ms < entry.cached_at_unix_ms + || !canonical_media_type(&entry.media_type) + || entry.extension.is_empty() + || validate_dimensions(Some(entry.width), Some(entry.height)).is_err() + }) + { + return Err(Phase1InboundMediaError::CorruptState); + } + if self + .configuration + .is_some_and(|value| value.validate().is_err()) + { + return Err(Phase1InboundMediaError::CorruptState); + } + self.total_bytes().map(|_| ()) + } +} + +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub struct Phase1MediaCacheStatus { + pub artifacts: u32, + pub bytes: u64, + pub configuration: Option<Phase1MediaConfigurationFingerprint>, +} + +impl Phase1MediaCacheIndex { + pub fn status(&self) -> Result<Phase1MediaCacheStatus, Phase1InboundMediaError> { + self.validate()?; + Ok(Phase1MediaCacheStatus { + artifacts: self.artifact_count(), + bytes: self.total_bytes()?, + configuration: self.configuration, + }) + } +} + +#[derive(Clone, Debug, Error, Eq, PartialEq)] +pub enum Phase1InboundMediaError { + #[error("inbound media reference is invalid")] + InvalidReference, + #[error("inbound media digest is invalid")] + InvalidDigest, + #[error("inbound media reference requires a digest")] + MissingDigest, + #[error("inbound media type is invalid")] + InvalidMediaType, + #[error("inbound media dimensions are invalid")] + InvalidDimensions, + #[error("inbound media byte size is invalid")] + InvalidByteSize, + #[error("inbound media alternative text is invalid")] + InvalidAlt, + #[error("inbound media metadata does not match verified bytes")] + MetadataMismatch, + #[error("inbound media operation is invalid")] + InvalidOperation, + #[error("inbound media operation identity does not match")] + OperationMismatch, + #[error("inbound media failure evidence is invalid")] + InvalidFailure, + #[error("inbound media configuration is invalid")] + InvalidConfiguration, + #[error("inbound media configuration changed")] + ConfigurationMismatch, + #[error("inbound media verification time is invalid")] + InvalidVerificationTime, + #[error("inbound media cache policy is invalid")] + InvalidCachePolicy, + #[error("inbound media cache observation is invalid")] + InvalidCacheObservation, + #[error("inbound media artifact exceeds cache quota")] + CacheQuotaExceeded, + #[error("inbound media artifact identity collides with different metadata")] + ArtifactCollision, + #[error("inbound media receipt is corrupt")] + CorruptReceipt, + #[error("inbound media state is corrupt")] + CorruptState, + #[error("inbound media schema version is unsupported")] + UnsupportedSchema, +} + +fn validate_url_text(value: &str) -> Result<(), Phase1InboundMediaError> { + if value.is_empty() + || value.len() > MEDIA_URL_MAX_BYTES + || value.chars().any(|character| character.is_control()) + { + return Err(Phase1InboundMediaError::InvalidReference); + } + Ok(()) +} + +fn canonical_media_type(value: &str) -> bool { + MediaType::parse(value).is_ok_and(|parsed| parsed.to_string() == value) +} + +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(()), + _ => Err(Phase1InboundMediaError::InvalidDimensions), + } +} + +fn update_optional(digest: &mut Sha256Hasher, value: Option<&str>) { + match value { + Some(value) => { + digest.update([1]); + digest.update((value.len() as u64).to_be_bytes()); + digest.update(value.as_bytes()); + } + None => digest.update([0]), + } +} + +fn update_optional_u64(digest: &mut Sha256Hasher, value: Option<u64>) { + match value { + Some(value) => { + digest.update([1]); + digest.update(value.to_be_bytes()); + } + None => digest.update([0]), + } +} + +#[cfg(test)] +mod tests { + use super::*; + + const HASH: &str = "2cf24dba5fb0a30e26e83b2ac5b9e29e1b161e5c1fa7425e73043362938b9824"; + + fn reference(alt: Option<&str>) -> Phase1StructuralMediaReference { + Phase1StructuralMediaReference::new( + format!("https://media.example/{HASH}.jpg"), + Some(HASH.to_owned()), + Some("image/jpeg".to_owned()), + Some(2), + Some(3), + Some(5), + alt.map(str::to_owned), + ) + .expect("reference") + } + + fn configuration(value: u8) -> Phase1MediaConfigurationFingerprint { + Phase1MediaConfigurationFingerprint::new([value; 32]).expect("configuration") + } + + fn receipt( + reference: &Phase1StructuralMediaReference, + configuration: Phase1MediaConfigurationFingerprint, + ) -> Phase1VerifiedMediaReceipt { + let commitment = + ByteCommitment::from_bytes(b"hello", MediaType::parse("image/jpeg").unwrap()); + Phase1VerifiedMediaReceipt::from_commitment( + reference, + BlobUrl::parse(&format!("https://cdn.example/{HASH}.jpg")).unwrap(), + &commitment, + 2, + 3, + configuration, + 10, + ) + .expect("receipt") + } + + #[test] + fn structural_reference_is_canonical_and_metadata_sensitive() { + let first = reference(Some("Harvest")); + let second = reference(Some("Harvest detail")); + assert_ne!(first.fingerprint(), second.fingerprint()); + assert_eq!(first.expected_sha256(), Some(HASH)); + assert!( + Phase1StructuralMediaReference::new( + format!("https://media.example/{}.jpg", "a".repeat(64)), + Some(HASH.to_owned()), + Some("image/jpeg".to_owned()), + Some(2), + Some(3), + Some(5), + None, + ) + .is_err() + ); + let interoperable = Phase1StructuralMediaReference::new( + "https://cdn.example/harvest.jpg", + Some(HASH.to_owned()), + Some("image/jpeg".to_owned()), + Some(2), + Some(3), + Some(5), + None, + ) + .expect("non-Blossom NIP-92 reference remains structural"); + assert!( + Phase1VerifiedMediaReceipt::from_commitment( + &interoperable, + BlobUrl::parse(&format!("https://cdn.example/{HASH}.jpg")).unwrap(), + &ByteCommitment::from_bytes(b"hello", MediaType::parse("image/jpeg").unwrap()), + 2, + 3, + configuration(1), + 1, + ) + .is_err() + ); + } + + #[test] + fn verified_state_requires_matching_operation_bytes_and_configuration() { + let structural = reference(None); + let mut media = MediaReference::new(structural.clone()).unwrap(); + let pending = Phase1InboundMediaPending::new([7; 16], configuration(3), 9).unwrap(); + media.begin(pending).unwrap(); + assert_eq!( + media.verify([8; 16], receipt(&structural, configuration(3))), + Err(Phase1InboundMediaError::OperationMismatch) + ); + media + .verify([7; 16], receipt(&structural, configuration(3))) + .unwrap(); + assert!(matches!( + media.retrieval(), + Phase1InboundMediaState::Verified(_) + )); + } + + #[test] + fn receipt_rejects_hash_size_type_and_dimension_mismatch() { + let expected = reference(None); + let wrong_bytes = + ByteCommitment::from_bytes(b"other", MediaType::parse("image/jpeg").unwrap()); + assert!( + Phase1VerifiedMediaReceipt::from_commitment( + &expected, + BlobUrl::parse(&format!("https://cdn.example/{}.jpg", wrong_bytes.sha256())) + .unwrap(), + &wrong_bytes, + 2, + 3, + configuration(1), + 1, + ) + .is_err() + ); + let commitment = + ByteCommitment::from_bytes(b"hello", MediaType::parse("image/png").unwrap()); + assert!( + Phase1VerifiedMediaReceipt::from_commitment( + &expected, + BlobUrl::parse(&format!("https://cdn.example/{HASH}.jpg")).unwrap(), + &commitment, + 2, + 3, + configuration(1), + 1, + ) + .is_err() + ); + } + + #[test] + fn cache_is_content_addressed_bounded_lru_and_configuration_scoped() { + let config = configuration(4); + let first_reference = reference(None); + let first = receipt(&first_reference, config); + let second_hash = Sha256::digest(b"world").to_hex(); + let second_reference = Phase1StructuralMediaReference::new( + format!("https://media.example/{second_hash}.jpg"), + Some(second_hash.clone()), + Some("image/jpeg".to_owned()), + Some(2), + Some(3), + Some(5), + None, + ) + .unwrap(); + let second_commitment = + ByteCommitment::from_bytes(b"world", MediaType::parse("image/jpeg").unwrap()); + let second = Phase1VerifiedMediaReceipt::from_commitment( + &second_reference, + BlobUrl::parse(&format!("https://cdn.example/{second_hash}.jpg")).unwrap(), + &second_commitment, + 2, + 3, + config, + 11, + ) + .unwrap(); + let mut cache = Phase1MediaCacheIndex::default(); + let policy = Phase1MediaCachePolicy::new(5, 1).unwrap(); + assert!(cache.admit(&first, policy, 10).unwrap().is_empty()); + let evicted = cache.admit(&second, policy, 11).unwrap(); + assert_eq!(evicted, vec![first.artifact_id()]); + assert!(!cache.contains(&first)); + assert!(cache.contains(&second)); + assert_eq!(cache.status().unwrap().artifacts, 1); + assert_eq!( + cache.admit(&second, policy, 12), + Ok(Vec::new()), + "idempotent cache admission remains bounded" + ); + assert_eq!( + cache.admit(&receipt(&first_reference, configuration(5)), policy, 13,), + Err(Phase1InboundMediaError::ConfigurationMismatch) + ); + assert_eq!( + cache.invalidate_configuration(configuration(5)), + vec![second.artifact_id()] + ); + assert_eq!(cache.status().unwrap().artifacts, 0); + } + + #[test] + fn persisted_receipt_and_cache_tamper_fail_closed() { + let structural = reference(None); + let config = configuration(7); + let mut media = MediaReference::new(structural.clone()).unwrap(); + media + .begin(Phase1InboundMediaPending::new([8; 16], config, 1).unwrap()) + .unwrap(); + let verified = receipt(&structural, config); + media.verify([8; 16], verified.clone()).unwrap(); + let mut media_value = serde_json::to_value(&media).unwrap(); + media_value["retrieval"]["evidence"]["observedSha256"] = serde_json::json!("a".repeat(64)); + let corrupt: MediaReference = serde_json::from_value(media_value).unwrap(); + assert_eq!( + corrupt.validate(), + Err(Phase1InboundMediaError::CorruptReceipt) + ); + + let mut cache = Phase1MediaCacheIndex::default(); + cache + .admit(&verified, Phase1MediaCachePolicy::new(10, 1).unwrap(), 12) + .unwrap(); + let mut cache_value = serde_json::to_value(cache).unwrap(); + let entry = cache_value["entries"] + .as_object_mut() + .unwrap() + .values_mut() + .next() + .unwrap(); + entry["byteSize"] = serde_json::json!(0); + let corrupt: Phase1MediaCacheIndex = serde_json::from_value(cache_value).unwrap(); + assert_eq!(corrupt.status(), Err(Phase1InboundMediaError::CorruptState)); + } +} diff --git a/crates/mobile_core/src/runtime/product_surface/model.rs b/crates/mobile_core/src/runtime/product_surface/model.rs @@ -1,6 +1,6 @@ use serde::{Deserialize, Serialize}; -use super::{CardId, ContextRank, TodayRank}; +use super::{CardId, ContextRank, MediaReference, TodayRank}; /// The closed Phase 1 top-level Today taxonomy. #[derive(Clone, Copy, Debug, Deserialize, Eq, Hash, Ord, PartialEq, PartialOrd, Serialize)] @@ -93,30 +93,6 @@ pub enum SupportingProfile { Deletion, } -/// Local media verification never changes the canonical card type. -#[derive(Clone, Copy, Debug, Deserialize, Eq, PartialEq, Serialize)] -#[serde(rename_all = "PascalCase")] -pub enum MediaVerificationState { - Pending, - Verified, - Failed, - Unavailable, -} - -/// Structural media metadata plus its separate local verification state. -#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)] -#[serde(rename_all = "camelCase")] -pub struct MediaReference { - pub url: String, - pub sha256: Option<String>, - pub media_type: Option<String>, - pub width: Option<u32>, - pub height: Option<u32>, - pub byte_size: Option<u64>, - pub alt: Option<String>, - pub verification: MediaVerificationState, -} - /// Tolerant profile attribution attached to cards and Me results. #[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)] #[serde(rename_all = "camelCase")] @@ -272,16 +248,4 @@ mod tests { r#"["Update","PhotoUpdate","Ask","Event","FoodAvailability"]"# ); } - - #[test] - fn media_state_is_independent_from_card_type() { - for state in [ - MediaVerificationState::Pending, - MediaVerificationState::Verified, - MediaVerificationState::Failed, - MediaVerificationState::Unavailable, - ] { - assert!(!serde_json::to_string(&state).expect("state").is_empty()); - } - } } diff --git a/crates/mobile_core/src/runtime/product_surface/projection.rs b/crates/mobile_core/src/runtime/product_surface/projection.rs @@ -6,7 +6,7 @@ use serde::{Deserialize, Serialize}; use super::{ CardId, CardLifecycleState, CardSourceIdentity, ClassifiedCard, ContextAdmission, - LocalNetworkAdmission, MediaReference, MediaVerificationState, SupportingProfile, + LocalNetworkAdmission, MediaReference, Phase1StructuralMediaReference, SupportingProfile, TodayCardType, }; @@ -254,16 +254,15 @@ fn post_media( .iter() .filter_map(|media| { let dimensions = media.dimensions(); - Some(MediaReference { - url: media.url()?.to_owned(), - sha256: media.sha256().map(str::to_owned), - media_type: media.media_type().map(str::to_owned), - width: dimensions.map(|value| value.width()), - height: dimensions.map(|value| value.height()), - byte_size: media.size(), - alt: media.alt().map(str::to_owned), - verification: MediaVerificationState::Unavailable, - }) + media_reference( + media.url()?, + media.sha256().map(str::to_owned), + media.media_type().map(str::to_owned), + dimensions.map(|value| value.width()), + dimensions.map(|value| value.height()), + media.size(), + media.alt().map(str::to_owned), + ) }) .collect() } @@ -276,16 +275,15 @@ fn food_media( .iter() .filter_map(|media| { let dimensions = media.dimensions(); - Some(MediaReference { - url: media.url()?.to_owned(), - sha256: blossom_digest(media.url()?), - media_type: None, - width: dimensions.map(|value| value.width()), - height: dimensions.map(|value| value.height()), - byte_size: None, - alt: None, - verification: MediaVerificationState::Unavailable, - }) + media_reference( + media.url()?, + blossom_digest(media.url()?), + None, + dimensions.map(|value| value.width()), + dimensions.map(|value| value.height()), + None, + None, + ) }) .collect() } @@ -294,20 +292,26 @@ fn calendar_media(tags: Vec<Vec<String>>) -> Vec<MediaReference> { tags.into_iter() .find(|tag| tag.first().map(String::as_str) == Some("image")) .and_then(|tag| tag.get(1).cloned()) - .map(|url| MediaReference { - sha256: blossom_digest(&url), - url, - media_type: None, - width: None, - height: None, - byte_size: None, - alt: None, - verification: MediaVerificationState::Unavailable, - }) + .and_then(|url| media_reference(&url, blossom_digest(&url), None, None, None, None, None)) .into_iter() .collect() } +#[allow(clippy::too_many_arguments)] +fn media_reference( + url: &str, + sha256: Option<String>, + media_type: Option<String>, + width: Option<u32>, + height: Option<u32>, + byte_size: Option<u64>, + alt: Option<String>, +) -> Option<MediaReference> { + Phase1StructuralMediaReference::new(url, sha256, media_type, width, height, byte_size, alt) + .and_then(MediaReference::new) + .ok() +} + fn blossom_digest(url: &str) -> Option<String> { let path = url.split_once("://")?.1.split_once('/')?.1; let candidate = path.split(['.', '/', '?', '#']).next()?; diff --git a/crates/mobile_core/src/runtime/product_surface/today.rs b/crates/mobile_core/src/runtime/product_surface/today.rs @@ -34,10 +34,14 @@ use radroots_transport::{ use super::{ CardId, CardLifecycleState, ClassifiedCard, CursorError, CursorScope, LocalAuthorOverlay, - LocalNetwork, LocalityEvidence, MeSnapshot, MediaReference, MediaVerificationState, - ProductEventClassification, ProfileSummary, SearchResult, SearchResultType, SupportingProfile, - ThreadEntry, ThreadReference, TimeRelevance, TodayCard, TodayCardType, TodayCursor, - TodayCursorPosition, TodayPage, TodayRank, TodayRankInput, classify_admitted_event, + LocalNetwork, LocalityEvidence, MeSnapshot, MediaReference, Phase1InboundMediaError, + Phase1InboundMediaFailure, Phase1InboundMediaPending, Phase1InboundMediaState, + Phase1MediaArtifactId, Phase1MediaCacheIndex, Phase1MediaCachePolicy, Phase1MediaCacheStatus, + Phase1MediaConfigurationFingerprint, Phase1StructuralMediaReference, + Phase1VerifiedMediaReceipt, ProductEventClassification, ProfileSummary, SearchResult, + SearchResultType, SupportingProfile, ThreadEntry, ThreadReference, TimeRelevance, TodayCard, + TodayCardType, TodayCursor, TodayCursorPosition, TodayPage, TodayRank, TodayRankInput, + classify_admitted_event, }; use crate::runtime::RadrootsRuntime; @@ -153,6 +157,8 @@ pub enum TodayError { Storage(#[from] radroots_storage::Error), #[error("today projection serialization failed")] Serialization, + #[error(transparent)] + InboundMedia(#[from] Phase1InboundMediaError), } #[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)] @@ -182,6 +188,8 @@ struct TodayProjectionState { profiles: BTreeMap<String, ProfileSummary>, thread: Vec<ThreadEntry>, overlays: BTreeMap<String, LocalAuthorOverlay>, + #[serde(default)] + media_cache: Phase1MediaCacheIndex, } #[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)] @@ -313,15 +321,7 @@ impl RadrootsRuntime { let generation = projection_generation()?; let projection_id = projection_id()?; let key = projection_document_key(context); - let prior = ProjectionStore::projection_document( - storage, - projection_id.clone(), - generation, - key.clone(), - ) - .await? - .map(|document| decode_state(document.value())) - .transpose()?; + let prior = load_state(storage, context, generation).await?; if update == TodayProjectionUpdate::Incremental && prior @@ -334,7 +334,11 @@ impl RadrootsRuntime { let visible = query_all_visible(storage).await?; let local_media = prior.as_ref().map(local_media_evidence).unwrap_or_default(); - let overlays = prior.map_or_else(BTreeMap::new, |state| state.overlays); + let overlays = prior + .as_ref() + .map_or_else(BTreeMap::new, |state| state.overlays.clone()); + let media_cache = + prior.map_or_else(Phase1MediaCacheIndex::default, |state| state.media_cache); let mut state = project_state( context, event_status.generation().as_bytes(), @@ -342,6 +346,7 @@ impl RadrootsRuntime { visible, overlays, )?; + state.media_cache = media_cache; apply_local_media_evidence(&mut state, &local_media); state.content_generation = content_generation(&state)?; let encoded = encode(&state)?; @@ -421,9 +426,13 @@ impl RadrootsRuntime { return Err(CursorError::Stale.into()); } let position = TodayCursor::decode(cursor, &scope)?; - let snapshot = load_snapshot(storage, projection_id, algorithm_generation, &scope) + let mut snapshot = load_snapshot(storage, projection_id, algorithm_generation, &scope) .await? .ok_or(TodayError::SnapshotMissing)?; + let current = load_state(storage, context, algorithm_generation) + .await? + .ok_or(TodayError::ProjectionMissing)?; + sanitize_snapshot_media(&mut snapshot, &current.media_cache); (scope, snapshot, Some(position.rank)) } else { let as_of = request @@ -558,16 +567,75 @@ impl RadrootsRuntime { }) } - /// Updates local byte-verification evidence without changing card classification. - pub async fn phase1_set_media_verification( + /// Starts one typed retrieval for every occurrence of the exact structural + /// reference. URL equality alone is deliberately insufficient. + pub async fn phase1_begin_media_retrieval( &self, context: &LocalNetwork, - url: &str, - verification: MediaVerificationState, + reference_fingerprint: [u8; 32], + pending: Phase1InboundMediaPending, ) -> Result<bool, TodayError> { - if url.is_empty() || url.len() > 8_192 || url.chars().any(char::is_control) { - return Err(TodayError::InvalidRequest); + 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 mut trial = state.clone(); + let prior_configuration = trial.media_cache.status()?.configuration; + if prior_configuration.is_some_and(|value| value != pending.configuration()) { + return Err(Phase1InboundMediaError::ConfigurationMismatch.into()); + } + trial + .media_cache + .invalidate_configuration(pending.configuration()); + let changed = mutate_matching_media(&mut trial, reference_fingerprint, |media| { + media.begin(pending.clone()) + })?; + if changed { + state = trial; + persist_media_state(storage, context, generation, &mut state).await?; + } + Ok(changed) + } + + /// Records a bounded, safe retrieval failure for the active operation. + pub async fn phase1_fail_media_retrieval( + &self, + context: &LocalNetwork, + reference_fingerprint: [u8; 32], + failure: Phase1InboundMediaFailure, + ) -> Result<bool, TodayError> { + 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 changed = mutate_matching_media(&mut state, reference_fingerprint, |media| { + media.fail(failure.clone()) + })?; + if changed { + persist_media_state(storage, context, generation, &mut state).await?; } + Ok(changed) + } + + /// 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( + &self, + context: &LocalNetwork, + reference_fingerprint: [u8; 32], + operation_id: [u8; 16], + receipt: Phase1VerifiedMediaReceipt, + policy: Phase1MediaCachePolicy, + cached_at_unix_ms: u64, + ) -> Result<Vec<Phase1MediaArtifactId>, TodayError> { let storage = self .client .storage() @@ -576,34 +644,108 @@ impl RadrootsRuntime { let mut state = load_state(storage, context, generation) .await? .ok_or(TodayError::ProjectionMissing)?; - let mut changed = false; - for projected in &mut state.cards { - for media in &mut projected.card.media { - if media.url == url && media.verification != verification { - media.verification = verification; - changed = true; - } - } + let mut trial = state.clone(); + let changed = mutate_matching_media(&mut trial, reference_fingerprint, |media| { + media.verify(operation_id, receipt.clone()) + })?; + if !changed { + return Err(TodayError::InvalidRequest); } - for profile in state.profiles.values_mut() { - for media in [&mut profile.picture, &mut profile.banner] - .into_iter() - .flatten() - { - if media.url == url && media.verification != verification { - media.verification = verification; - changed = true; - } - } + let evicted = trial + .media_cache + .admit(&receipt, policy, cached_at_unix_ms)?; + for artifact_id in &evicted { + invalidate_artifact_references(&mut trial, *artifact_id); } + state = trial; + persist_media_state(storage, context, generation, &mut state).await?; + Ok(evicted) + } + + /// Records one successful local artifact access for deterministic LRU. + pub async fn phase1_touch_media_artifact( + &self, + context: &LocalNetwork, + artifact_id: Phase1MediaArtifactId, + observed_at_unix_ms: u64, + ) -> Result<bool, TodayError> { + 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 changed = state.media_cache.touch(artifact_id, observed_at_unix_ms)?; if changed { - refresh_thread_profiles(&mut state); - state.content_generation = content_generation(&state)?; - store_state(storage, context, generation, &state).await?; + persist_media_state(storage, context, generation, &mut state).await?; } Ok(changed) } + /// Invalidates a missing, corrupt, or explicitly evicted local artifact. + pub async fn phase1_invalidate_media_artifact( + &self, + context: &LocalNetwork, + artifact_id: Phase1MediaArtifactId, + ) -> Result<bool, TodayError> { + 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 cache_changed = state.media_cache.invalidate_artifact(artifact_id); + let references_changed = invalidate_artifact_references(&mut state, artifact_id); + if cache_changed || references_changed { + persist_media_state(storage, context, generation, &mut state).await?; + } + Ok(cache_changed || references_changed) + } + + /// Clears all trust derived under an obsolete endpoint/network/cache + /// configuration and records the new configuration generation. + pub async fn phase1_invalidate_media_configuration( + &self, + context: &LocalNetwork, + configuration: Phase1MediaConfigurationFingerprint, + ) -> Result<Vec<Phase1MediaArtifactId>, TodayError> { + 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 prior_configuration = state.media_cache.status()?.configuration; + let removed = state.media_cache.invalidate_configuration(configuration); + if prior_configuration != Some(configuration) { + for_each_media_mut(&mut state, |media| { + media.invalidate(); + }); + persist_media_state(storage, context, generation, &mut state).await?; + } + Ok(removed) + } + + pub async fn phase1_media_cache_status( + &self, + context: &LocalNetwork, + ) -> Result<Phase1MediaCacheStatus, TodayError> { + let storage = self + .client + .storage() + .map_err(|_| TodayError::RuntimeUnavailable)?; + let state = load_state(storage, context, projection_generation()?) + .await? + .ok_or(TodayError::ProjectionMissing)?; + state.media_cache.status().map_err(TodayError::from) + } + /// Persists active-author delivery state as a local-only Today overlay. pub async fn phase1_set_local_author_overlay( &self, @@ -709,18 +851,20 @@ fn refresh_receipt( } } -fn local_media_evidence(state: &TodayProjectionState) -> BTreeMap<String, MediaVerificationState> { +fn local_media_evidence( + state: &TodayProjectionState, +) -> BTreeMap<[u8; 32], Phase1InboundMediaState> { let mut evidence = state .cards .iter() .flat_map(|projected| projected.card.media.iter()) - .filter(|media| media.verification != MediaVerificationState::Unavailable) - .map(|media| (media.url.clone(), media.verification)) + .filter(|media| !matches!(media.retrieval(), Phase1InboundMediaState::Unavailable)) + .map(|media| (*media.structural().fingerprint(), media.retrieval().clone())) .collect::<BTreeMap<_, _>>(); for profile in state.profiles.values() { for media in [&profile.picture, &profile.banner].into_iter().flatten() { - if media.verification != MediaVerificationState::Unavailable { - evidence.insert(media.url.clone(), media.verification); + if !matches!(media.retrieval(), Phase1InboundMediaState::Unavailable) { + evidence.insert(*media.structural().fingerprint(), media.retrieval().clone()); } } } @@ -729,13 +873,26 @@ fn local_media_evidence(state: &TodayProjectionState) -> BTreeMap<String, MediaV fn apply_local_media_evidence( state: &mut TodayProjectionState, - evidence: &BTreeMap<String, MediaVerificationState>, + evidence: &BTreeMap<[u8; 32], Phase1InboundMediaState>, +) { + let cache = state.media_cache.clone(); + for_each_media_mut(state, |media| { + if let Some(retrieval) = evidence.get(media.structural().fingerprint()) + && media.restore(retrieval.clone(), &cache).is_err() + { + media.invalidate(); + } + }); + refresh_thread_profiles(state); +} + +fn for_each_media_mut( + state: &mut TodayProjectionState, + mut action: impl FnMut(&mut MediaReference), ) { for projected in &mut state.cards { for media in &mut projected.card.media { - if let Some(verification) = evidence.get(&media.url) { - media.verification = *verification; - } + action(media); } } for profile in state.profiles.values_mut() { @@ -743,12 +900,65 @@ fn apply_local_media_evidence( .into_iter() .flatten() { - if let Some(verification) = evidence.get(&media.url) { - media.verification = *verification; + action(media); + } + } +} + +fn mutate_matching_media( + state: &mut TodayProjectionState, + reference_fingerprint: [u8; 32], + mut action: impl FnMut(&mut MediaReference) -> Result<(), Phase1InboundMediaError>, +) -> Result<bool, TodayError> { + let mut found = false; + let mut failure = None; + for_each_media_mut(state, |media| { + if failure.is_none() && media.structural().fingerprint() == &reference_fingerprint { + found = true; + if let Err(error) = action(media) { + failure = Some(error); } } + }); + if let Some(error) = failure { + return Err(error.into()); + } + if found { + refresh_thread_profiles(state); + } + Ok(found) +} + +fn invalidate_artifact_references( + state: &mut TodayProjectionState, + artifact_id: Phase1MediaArtifactId, +) -> bool { + let mut changed = false; + for_each_media_mut(state, |media| { + if matches!( + media.retrieval(), + Phase1InboundMediaState::Verified(receipt) if receipt.artifact_id() == artifact_id + ) { + media.invalidate(); + changed = true; + } + }); + if changed { + refresh_thread_profiles(state); } + changed +} + +async fn persist_media_state( + storage: &dyn radroots_storage::Storage, + context: &LocalNetwork, + generation: ProjectionGeneration, + state: &mut TodayProjectionState, +) -> Result<(), TodayError> { refresh_thread_profiles(state); + validate_media_state(state)?; + state.content_generation = content_generation(state)?; + store_state(storage, context, generation, state).await } fn refresh_thread_profiles(state: &mut TodayProjectionState) { @@ -836,6 +1046,7 @@ fn project_state( profiles, thread, overlays, + media_cache: Phase1MediaCacheIndex::default(), }) } @@ -851,10 +1062,12 @@ fn profile_summary(admitted: &RadrootsAdmittedEvent) -> Result<ProfileSummary, T about: metadata.about().map(str::to_owned), picture: metadata .picture() - .map(|value| unverified_media(value.as_str())), + .map(|value| unverified_media(value.as_str())) + .transpose()?, banner: metadata .banner() - .map(|value| unverified_media(value.as_str())), + .map(|value| unverified_media(value.as_str())) + .transpose()?, nip05: metadata.nip05().map(|value| value.as_str().to_owned()), website: typed_profile_extra(metadata.raw_fields(), "website"), lightning_address: typed_profile_extra(metadata.raw_fields(), "lud16"), @@ -871,17 +1084,17 @@ fn typed_profile_extra(fields: &BTreeMap<String, serde_json::Value>, key: &str) .map(str::to_owned) } -fn unverified_media(url: &str) -> MediaReference { - MediaReference { - url: url.to_owned(), - sha256: blossom_digest(url), - media_type: None, - width: None, - height: None, - byte_size: None, - alt: None, - verification: MediaVerificationState::Unavailable, - } +fn unverified_media(url: &str) -> Result<MediaReference, TodayError> { + MediaReference::new(Phase1StructuralMediaReference::new( + url, + blossom_digest(url), + None, + None, + None, + None, + None, + )?) + .map_err(TodayError::from) } fn reply_entry(admitted: &RadrootsAdmittedEvent) -> Result<ThreadEntry, TodayError> { @@ -1098,15 +1311,21 @@ async fn load_state( context: &LocalNetwork, generation: ProjectionGeneration, ) -> Result<Option<TodayProjectionState>, TodayError> { - ProjectionStore::projection_document( + let document = ProjectionStore::projection_document( storage, projection_id()?, generation, projection_document_key(context), ) - .await? - .map(|document| decode_state(document.value())) - .transpose() + .await?; + let Some(document) = document else { + return Ok(None); + }; + let (state, migrated) = decode_state_document(document.value())?; + if migrated { + store_state(storage, context, generation, &state).await?; + } + Ok(Some(state)) } async fn store_state( @@ -1165,25 +1384,237 @@ async fn load_snapshot( .transpose() } +#[cfg(test)] fn decode_state(value: &[u8]) -> Result<TodayProjectionState, TodayError> { - let state: TodayProjectionState = decode(value)?; + decode_state_document(value).map(|(state, _)| state) +} + +fn decode_state_document(value: &[u8]) -> Result<(TodayProjectionState, bool), TodayError> { + if let Ok(state) = serde_json::from_slice::<TodayProjectionState>(value) + && state.schema_version == TODAY_PROJECTION_DOCUMENT_SCHEMA_VERSION + && state.content_generation != 0 + && content_generation(&state)? == state.content_generation + && validate_media_state(&state).is_ok() + { + return Ok((state, false)); + } + let state = migrate_legacy_state(value)?; if state.schema_version != TODAY_PROJECTION_DOCUMENT_SCHEMA_VERSION || state.content_generation == 0 || content_generation(&state)? != state.content_generation + || validate_media_state(&state).is_err() { return Err(TodayError::CorruptProjection); } - Ok(state) + Ok((state, true)) +} + +fn validate_media_state(state: &TodayProjectionState) -> Result<(), Phase1InboundMediaError> { + state.media_cache.status()?; + for projected in &state.cards { + for media in &projected.card.media { + media.validate()?; + if let Phase1InboundMediaState::Verified(receipt) = media.retrieval() + && !state.media_cache.contains(receipt) + { + return Err(Phase1InboundMediaError::CorruptState); + } + } + } + for profile in state.profiles.values() { + for media in [&profile.picture, &profile.banner].into_iter().flatten() { + media.validate()?; + if let Phase1InboundMediaState::Verified(receipt) = media.retrieval() + && !state.media_cache.contains(receipt) + { + return Err(Phase1InboundMediaError::CorruptState); + } + } + } + if state + .thread + .iter() + .any(|entry| entry.author_profile.as_ref() != state.profiles.get(&entry.author_pubkey)) + { + return Err(Phase1InboundMediaError::CorruptState); + } + Ok(()) +} + +fn sanitize_snapshot_media(snapshot: &mut FrozenTodaySnapshot, cache: &Phase1MediaCacheIndex) { + for item in &mut snapshot.items { + for media in &mut item.card.media { + if let Phase1InboundMediaState::Verified(receipt) = media.retrieval() + && !cache.contains(receipt) + { + media.invalidate(); + } + } + if let Some(profile) = &mut item.author_profile { + for media in [&mut profile.picture, &mut profile.banner] + .into_iter() + .flatten() + { + if let Phase1InboundMediaState::Verified(receipt) = media.retrieval() + && !cache.contains(receipt) + { + media.invalidate(); + } + } + } + for entry in &mut item.thread { + if let Some(profile) = &mut entry.author_profile { + for media in [&mut profile.picture, &mut profile.banner] + .into_iter() + .flatten() + { + if let Phase1InboundMediaState::Verified(receipt) = media.retrieval() + && !cache.contains(receipt) + { + media.invalidate(); + } + } + } + } + } } fn decode_snapshot(value: &[u8]) -> Result<FrozenTodaySnapshot, TodayError> { - let snapshot: FrozenTodaySnapshot = decode(value)?; + let snapshot: FrozenTodaySnapshot = match serde_json::from_slice(value) { + Ok(snapshot) => snapshot, + Err(_) => { + let mut legacy: serde_json::Value = decode(value)?; + migrate_legacy_media_values(&mut legacy)?; + serde_json::from_value(legacy).map_err(|_| TodayError::CorruptProjection)? + } + }; if snapshot.schema_version != TODAY_SNAPSHOT_SCHEMA_VERSION { return Err(TodayError::CorruptProjection); } Ok(snapshot) } +#[derive(Deserialize)] +#[serde(deny_unknown_fields, rename_all = "camelCase")] +struct LegacyMediaReference { + url: String, + sha256: Option<String>, + media_type: Option<String>, + width: Option<u32>, + height: Option<u32>, + byte_size: Option<u64>, + alt: Option<String>, + verification: LegacyMediaVerificationState, +} + +#[derive(Deserialize)] +#[serde(rename_all = "PascalCase")] +enum LegacyMediaVerificationState { + Pending, + Verified, + Failed, + Unavailable, +} + +fn migrate_legacy_state(value: &[u8]) -> Result<TodayProjectionState, TodayError> { + verify_legacy_content_generation(value)?; + let mut legacy: serde_json::Value = decode(value)?; + migrate_legacy_media_values(&mut legacy)?; + let object = legacy + .as_object_mut() + .ok_or(TodayError::CorruptProjection)?; + if object.contains_key("mediaCache") { + return Err(TodayError::CorruptProjection); + } + object.insert( + "mediaCache".to_owned(), + serde_json::to_value(Phase1MediaCacheIndex::default()) + .map_err(|_| TodayError::Serialization)?, + ); + object.insert("contentGeneration".to_owned(), serde_json::json!(0)); + let mut state: TodayProjectionState = + serde_json::from_value(legacy).map_err(|_| TodayError::CorruptProjection)?; + state.content_generation = content_generation(&state)?; + Ok(state) +} + +fn verify_legacy_content_generation(value: &[u8]) -> Result<(), TodayError> { + let parsed: serde_json::Value = decode(value)?; + let expected = parsed + .get("contentGeneration") + .and_then(serde_json::Value::as_u64) + .filter(|value| *value != 0) + .ok_or(TodayError::CorruptProjection)?; + if parsed + .get("schemaVersion") + .and_then(serde_json::Value::as_u64) + != Some(u64::from(TODAY_PROJECTION_DOCUMENT_SCHEMA_VERSION)) + { + return Err(TodayError::CorruptProjection); + } + let marker = b"\"contentGeneration\":"; + let starts = value + .windows(marker.len()) + .enumerate() + .filter_map(|(index, bytes)| (bytes == marker).then_some(index)) + .collect::<Vec<_>>(); + let [start] = starts.as_slice() else { + return Err(TodayError::CorruptProjection); + }; + let number_start = start + marker.len(); + let number_end = value[number_start..] + .iter() + .position(|byte| !byte.is_ascii_digit()) + .map(|offset| number_start + offset) + .ok_or(TodayError::CorruptProjection)?; + if number_end == number_start { + return Err(TodayError::CorruptProjection); + } + let mut canonical = Vec::with_capacity(value.len()); + canonical.extend_from_slice(&value[..number_start]); + canonical.push(b'0'); + canonical.extend_from_slice(&value[number_end..]); + let digest = Sha256::digest([PROJECTION_CONTENT_DOMAIN, canonical.as_slice()].concat()); + let observed = u64::from_be_bytes(digest[..8].try_into().expect("digest prefix")).max(1); + (observed == expected) + .then_some(()) + .ok_or(TodayError::CorruptProjection) +} + +fn migrate_legacy_media_values(value: &mut serde_json::Value) -> Result<(), TodayError> { + match value { + serde_json::Value::Array(values) => { + for value in values { + migrate_legacy_media_values(value)?; + } + } + serde_json::Value::Object(object) + if object.contains_key("verification") && object.contains_key("url") => + { + let legacy: LegacyMediaReference = + serde_json::from_value(value.clone()).map_err(|_| TodayError::CorruptProjection)?; + let _legacy_verification = legacy.verification; + let migrated = MediaReference::legacy_unavailable( + legacy.url, + legacy.sha256, + legacy.media_type, + legacy.width, + legacy.height, + legacy.byte_size, + legacy.alt, + )?; + *value = serde_json::to_value(migrated).map_err(|_| TodayError::Serialization)?; + } + serde_json::Value::Object(object) => { + for value in object.values_mut() { + migrate_legacy_media_values(value)?; + } + } + _ => {} + } + Ok(()) +} + fn content_generation(state: &TodayProjectionState) -> Result<u64, TodayError> { let mut canonical = state.clone(); canonical.content_generation = 0; @@ -1263,6 +1694,9 @@ mod tests { use nostr::secp256k1::Message; use nostr::{Keys, SECP256K1}; + use radroots_blossom::{ + BlobUrl, MediaType, Sha256 as BlossomSha256, descriptor::ByteCommitment, + }; use radroots_event::{ SignedEvent, admission::{AdmissionPolicy, RawEvent, VisibilityPolicy}, @@ -1480,6 +1914,117 @@ mod tests { .expect("ingest") } + #[allow(clippy::too_many_arguments)] + async fn verify_inbound_media( + runtime: &RadrootsRuntime, + context: &LocalNetwork, + source_url: &str, + bytes: &[u8], + media_type: &str, + width: u32, + height: u32, + operation_id: [u8; 16], + ) { + let storage = runtime.client.storage().expect("storage"); + let state = load_state(storage, context, projection_generation().unwrap()) + .await + .unwrap() + .expect("projection"); + let reference = state + .cards + .iter() + .flat_map(|value| value.card.media.iter()) + .chain( + state + .profiles + .values() + .flat_map(|profile| [&profile.picture, &profile.banner].into_iter().flatten()), + ) + .find(|media| media.structural().source_url() == source_url) + .expect("structural media") + .structural() + .clone(); + let configuration = Phase1MediaConfigurationFingerprint::new([9; 32]).unwrap(); + runtime + .phase1_begin_media_retrieval( + context, + *reference.fingerprint(), + Phase1InboundMediaPending::new(operation_id, configuration, 10).unwrap(), + ) + .await + .expect("begin retrieval"); + let commitment = ByteCommitment::from_bytes(bytes, MediaType::parse(media_type).unwrap()); + let receipt = Phase1VerifiedMediaReceipt::from_commitment( + &reference, + BlobUrl::parse(source_url).unwrap(), + &commitment, + width, + height, + configuration, + 11, + ) + .expect("byte receipt"); + runtime + .phase1_commit_media_receipt( + context, + *reference.fingerprint(), + operation_id, + receipt, + Phase1MediaCachePolicy::new(1_024, 8).unwrap(), + 12, + ) + .await + .expect("commit receipt"); + } + + fn downgrade_media_to_legacy(value: &mut serde_json::Value) { + match value { + serde_json::Value::Array(values) => { + for value in values { + downgrade_media_to_legacy(value); + } + } + serde_json::Value::Object(object) + if object.contains_key("structural") && object.contains_key("retrieval") => + { + let structural = object + .get("structural") + .and_then(serde_json::Value::as_object) + .expect("structural object"); + *value = serde_json::json!({ + "url": structural.get("sourceUrl").cloned().unwrap(), + "sha256": structural.get("expectedSha256").cloned().unwrap(), + "mediaType": structural.get("expectedMediaType").cloned().unwrap(), + "width": structural.get("expectedWidth").cloned().unwrap(), + "height": structural.get("expectedHeight").cloned().unwrap(), + "byteSize": structural.get("expectedByteSize").cloned().unwrap(), + "alt": structural.get("alt").cloned().unwrap(), + "verification": "Verified", + }); + } + serde_json::Value::Object(object) => { + for value in object.values_mut() { + downgrade_media_to_legacy(value); + } + } + _ => {} + } + } + + fn legacy_state_bytes(state: &TodayProjectionState) -> Vec<u8> { + let mut value = serde_json::to_value(state).expect("state value"); + let object = value.as_object_mut().expect("state object"); + object.remove("mediaCache").expect("new cache field"); + object.insert("contentGeneration".to_owned(), serde_json::json!(0)); + downgrade_media_to_legacy(&mut value); + let canonical = serde_json::to_vec(&value).expect("legacy canonical"); + let digest = + sha2::Sha256::digest([PROJECTION_CONTENT_DOMAIN, canonical.as_slice()].concat()); + let generation = u64::from_be_bytes(digest[..8].try_into().expect("digest prefix")).max(1); + value["contentGeneration"] = serde_json::json!(generation); + serde_json::to_vec(&value).expect("legacy state") + } + #[cfg(feature = "mobile-social")] #[tokio::test] async fn relay_sync_fetches_the_exact_today_selector_and_projects_real_events() { @@ -1587,7 +2132,9 @@ mod tests { let runtime = RadrootsRuntime::test_memory().expect("runtime"); let context = context(Some("victoria"), 7); let author = keys().public_key().to_string(); - let profile_url = "https://blob.example/aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa.jpg"; + let profile_bytes = b"profile picture"; + let profile_digest = BlossomSha256::digest(profile_bytes).to_hex(); + let profile_url = format!("https://blob.example/{profile_digest}.jpg"); let profile = signed( 0, Vec::new(), @@ -1598,7 +2145,8 @@ mod tests { ); ingest(&runtime, &context, profile, 2_000_000_100).await; - let digest = "b".repeat(64); + let photo_bytes = b"photo"; + let digest = BlossomSha256::digest(photo_bytes).to_hex(); let photo_url = format!("https://blob.example/{digest}.jpg"); let photo_content = format!("Fresh field photo {photo_url}"); let photo = signed_owned( @@ -1611,7 +2159,7 @@ mod tests { format!("x {digest}"), "m image/jpeg".into(), "dim 120x80".into(), - "size 345".into(), + format!("size {}", photo_bytes.len()), "alt Fresh field photo".into(), ], ], @@ -1688,26 +2236,28 @@ mod tests { ); ingest(&runtime, &context, comment, 2_000_000_105).await; - assert!( - runtime - .phase1_set_media_verification( - &context, - &photo_url, - MediaVerificationState::Verified, - ) - .await - .expect("media evidence") - ); - assert!( - runtime - .phase1_set_media_verification( - &context, - profile_url, - MediaVerificationState::Verified, - ) - .await - .expect("profile media evidence") - ); + verify_inbound_media( + &runtime, + &context, + &photo_url, + photo_bytes, + "image/jpeg", + 120, + 80, + [1; 16], + ) + .await; + verify_inbound_media( + &runtime, + &context, + &profile_url, + profile_bytes, + "image/jpeg", + 24, + 24, + [2; 16], + ) + .await; let live = runtime .phase1_today_page(&context, TodayPageRequest::first(100, 2_000_000_200)) .await @@ -1719,10 +2269,10 @@ mod tests { .find(|card| card.card.card_type == TodayCardType::PhotoUpdate) .unwrap_or_else(|| panic!("photo card missing from {:#?}", live.items)); assert_eq!(card.card.card_type, TodayCardType::PhotoUpdate); - assert_eq!( - card.card.media[0].verification, - MediaVerificationState::Verified - ); + assert!(matches!( + card.card.media[0].retrieval(), + Phase1InboundMediaState::Verified(_) + )); assert_eq!(card.thread.len(), 1); assert_eq!( live.items @@ -1739,14 +2289,14 @@ mod tests { profile.lightning_address.as_deref(), Some("moss@example.com") ); - assert_eq!( + assert!(matches!( profile .picture .as_ref() .expect("profile picture") - .verification, - MediaVerificationState::Verified - ); + .retrieval(), + Phase1InboundMediaState::Verified(_) + )); runtime .phase1_set_local_author_overlay( @@ -1799,19 +2349,19 @@ mod tests { .iter() .find(|card| card.card.card_type == TodayCardType::PhotoUpdate) .expect("rebuilt photo"); - assert_eq!( - rebuilt_photo.card.media[0].verification, - MediaVerificationState::Verified - ); - assert_eq!( + assert!(matches!( + rebuilt_photo.card.media[0].retrieval(), + Phase1InboundMediaState::Verified(_) + )); + assert!(matches!( rebuilt_photo .author_profile .as_ref() .and_then(|profile| profile.picture.as_ref()) .expect("rebuilt profile picture") - .verification, - MediaVerificationState::Verified - ); + .retrieval(), + Phase1InboundMediaState::Verified(_) + )); assert_eq!( rebuilt_photo .local_overlay @@ -1820,6 +2370,176 @@ mod tests { Some("delivered") ); assert_eq!(rebuilt_photo.thread.len(), 1); + + let photo_reference = rebuilt_photo.card.media[0].structural().clone(); + let photo_receipt = match rebuilt_photo.card.media[0].retrieval() { + Phase1InboundMediaState::Verified(receipt) => (**receipt).clone(), + state => panic!("unexpected photo state: {state:?}"), + }; + let profile_artifact = match rebuilt_photo + .author_profile + .as_ref() + .and_then(|profile| profile.picture.as_ref()) + .expect("profile picture before eviction") + .retrieval() + { + Phase1InboundMediaState::Verified(receipt) => receipt.artifact_id(), + state => panic!("unexpected profile state: {state:?}"), + }; + runtime + .phase1_begin_media_retrieval( + &context, + *photo_reference.fingerprint(), + Phase1InboundMediaPending::new( + [3; 16], + Phase1MediaConfigurationFingerprint::new([9; 32]).unwrap(), + 30, + ) + .unwrap(), + ) + .await + .expect("repeat retrieval"); + let evicted = runtime + .phase1_commit_media_receipt( + &context, + *photo_reference.fingerprint(), + [3; 16], + photo_receipt, + Phase1MediaCachePolicy::new(1_024, 1).unwrap(), + 31, + ) + .await + .expect("quota commit"); + assert_eq!(evicted, vec![profile_artifact]); + let quota_page = runtime + .phase1_today_page(&context, TodayPageRequest::first(100, 2_000_000_203)) + .await + .expect("quota page"); + let quota_photo = quota_page + .items + .iter() + .find(|card| card.card.card_type == TodayCardType::PhotoUpdate) + .expect("quota photo"); + assert!(matches!( + quota_photo.card.media[0].retrieval(), + Phase1InboundMediaState::Verified(_) + )); + assert!(matches!( + quota_photo + .author_profile + .as_ref() + .and_then(|profile| profile.picture.as_ref()) + .expect("evicted profile picture") + .retrieval(), + Phase1InboundMediaState::Unavailable + )); + + let removed = runtime + .phase1_invalidate_media_configuration( + &context, + Phase1MediaConfigurationFingerprint::new([8; 32]).unwrap(), + ) + .await + .expect("configuration invalidation"); + assert_eq!(removed.len(), 1); + let invalidated = runtime + .phase1_today_page(&context, TodayPageRequest::first(100, 2_000_000_203)) + .await + .expect("page after configuration change"); + let invalidated_photo = invalidated + .items + .iter() + .find(|card| card.card.card_type == TodayCardType::PhotoUpdate) + .expect("invalidated photo"); + assert!(matches!( + invalidated_photo.card.media[0].retrieval(), + Phase1InboundMediaState::Unavailable + )); + assert!(matches!( + invalidated_photo + .author_profile + .as_ref() + .and_then(|profile| profile.picture.as_ref()) + .expect("invalidated profile picture") + .retrieval(), + Phase1InboundMediaState::Unavailable + )); + } + + #[tokio::test] + async fn legacy_enum_only_media_migrates_to_unavailable_and_repersists() { + let runtime = RadrootsRuntime::test_memory().expect("runtime"); + let context = context(None, 1); + let bytes = b"legacy profile"; + let hash = BlossomSha256::digest(bytes).to_hex(); + let url = format!("https://blob.example/{hash}.jpg"); + ingest( + &runtime, + &context, + signed( + 0, + Vec::new(), + &format!(r#"{{"name":"legacy","picture":"{url}"}}"#), + 2_000_000_000, + ), + 2_000_000_100, + ) + .await; + let storage = runtime.client.storage().expect("storage"); + let generation = projection_generation().unwrap(); + let state = load_state(storage, &context, generation) + .await + .unwrap() + .expect("state"); + let legacy = legacy_state_bytes(&state); + ProjectionStore::put_projection_document( + storage, + projection_id().unwrap(), + generation, + ProjectionDocument::new(projection_document_key(&context), legacy.clone()).unwrap(), + ) + .await + .unwrap(); + + let migrated = load_state(storage, &context, generation) + .await + .unwrap() + .expect("migrated state"); + let picture = migrated + .profiles + .values() + .next() + .and_then(|profile| profile.picture.as_ref()) + .expect("picture"); + assert!(matches!( + picture.retrieval(), + Phase1InboundMediaState::Unavailable + )); + assert_eq!(migrated.media_cache.status().unwrap().artifacts, 0); + let stored = ProjectionStore::projection_document( + storage, + projection_id().unwrap(), + generation, + projection_document_key(&context), + ) + .await + .unwrap() + .expect("repersisted state"); + let text = std::str::from_utf8(stored.value()).unwrap(); + assert!(!text.contains("\"verification\"")); + assert!(text.contains("\"mediaCache\"")); + + let mut tampered: serde_json::Value = serde_json::from_slice(&legacy).unwrap(); + tampered["profiles"] + .as_object_mut() + .unwrap() + .values_mut() + .next() + .unwrap()["name"] = serde_json::json!("tampered"); + assert!(matches!( + decode_state(&serde_json::to_vec(&tampered).unwrap()), + Err(TodayError::CorruptProjection) + )); } #[tokio::test] @@ -1936,6 +2656,104 @@ mod tests { } #[tokio::test] + async fn sqlite_reopen_preserves_receipts_and_missing_artifacts_fail_closed() { + use crate::runtime::{ + builder::RuntimeBuilder, + store::{MobileUserStoreConfig, ProtectedDataAvailability}, + }; + + let root = tempfile::tempdir().expect("root"); + let public_key = keys().public_key().to_string(); + let store = MobileUserStoreConfig::from_encoded( + root.path(), + &public_key, + "4141414141414141414141414141414141414141414141414141414141414141", + 2_000_000_000_000, + ProtectedDataAvailability::Available, + ) + .expect("store"); + std::fs::create_dir_all(store.owner_directory()).expect("owner directory"); + let context = context(None, 1); + let runtime = RuntimeBuilder::new(store.clone()) + .build() + .await + .expect("runtime"); + let bytes = b"sqlite profile"; + let hash = BlossomSha256::digest(bytes).to_hex(); + let url = format!("https://blob.example/{hash}.jpg"); + ingest( + &runtime, + &context, + signed( + 0, + Vec::new(), + &format!(r#"{{"name":"sqlite","picture":"{url}"}}"#), + 2_000_000_000, + ), + 2_000_000_100, + ) + .await; + verify_inbound_media( + &runtime, + &context, + &url, + bytes, + "image/jpeg", + 20, + 20, + [6; 16], + ) + .await; + assert_eq!( + runtime + .phase1_media_cache_status(&context) + .await + .unwrap() + .artifacts, + 1 + ); + runtime.shutdown().await.expect("shutdown"); + + let reopened = RuntimeBuilder::new(store).build().await.expect("reopen"); + let profile = reopened + .phase1_me(&context, &public_key, 2_000_000_200) + .await + .unwrap() + .profile + .expect("profile"); + let picture = profile.picture.expect("picture"); + let artifact_id = match picture.retrieval() { + Phase1InboundMediaState::Verified(receipt) => receipt.artifact_id(), + state => panic!("unexpected state: {state:?}"), + }; + assert!( + reopened + .phase1_invalidate_media_artifact(&context, artifact_id) + .await + .unwrap() + ); + assert_eq!( + reopened + .phase1_media_cache_status(&context) + .await + .unwrap() + .artifacts, + 0 + ); + let profile = reopened + .phase1_me(&context, &public_key, 2_000_000_201) + .await + .unwrap() + .profile + .expect("profile after invalidation"); + assert!(matches!( + profile.picture.unwrap().retrieval(), + Phase1InboundMediaState::Unavailable + )); + reopened.shutdown().await.expect("shutdown reopened"); + } + + #[tokio::test] async fn fail_closed_requests_cursors_overlays_and_projection_guards_are_executable() { let mut runtime = RadrootsRuntime::test_memory().expect("runtime"); let context = context(None, 1); @@ -2010,20 +2828,17 @@ mod tests { .await, Err(TodayError::InvalidRequest) )); - for url in ["", &"x".repeat(8_193), "bad\nurl"] { - assert!(matches!( - runtime - .phase1_set_media_verification(&context, url, MediaVerificationState::Verified,) - .await, - Err(TodayError::InvalidRequest) - )); - } assert!( !runtime - .phase1_set_media_verification( + .phase1_begin_media_retrieval( &context, - "https://blob.example/missing.jpg", - MediaVerificationState::Verified, + [3; 32], + Phase1InboundMediaPending::new( + [4; 16], + Phase1MediaConfigurationFingerprint::new([5; 32]).unwrap(), + 1, + ) + .unwrap(), ) .await .expect("unmatched media") diff --git a/crates/mobile_ffi/src/dto.rs b/crates/mobile_ffi/src/dto.rs @@ -22,9 +22,9 @@ use radroots_mobile_core::runtime::{ product_surface::{ AddCommandType, CardLifecycleState, CreateAsk, CreateEvent, CreateFoodAvailability, CreatePhotoUpdate, CreateUpdate, LocalNetwork, LocalNetworkRelayPolicy, MeSnapshot, - MediaReference, MediaVerificationState, Phase1AddCommand, Phase1CancellationPolicy, - Phase1DraftEventTiming, Phase1DraftFormSnapshot, Phase1DraftKind, Phase1DraftMediaSnapshot, - Phase1DraftStatus, Phase1MediaPrerequisite, Phase1MediaStage, Phase1OutboxState, + MediaReference, Phase1AddCommand, Phase1CancellationPolicy, Phase1DraftEventTiming, + Phase1DraftFormSnapshot, Phase1DraftKind, Phase1DraftMediaSnapshot, Phase1DraftStatus, + Phase1InboundMediaState, Phase1MediaPrerequisite, Phase1MediaStage, Phase1OutboxState, Phase1QueuePolicy, Phase1RelaySatisfaction, Phase1UploadIntent, ProfileSummary, SearchResult, SearchResultType, SupportingProfile, ThreadEntry, TodayCard, TodayCardType, TodayPage, TodayProjectionUpdate, TodayRefreshReceipt, TodayRelaySyncState, @@ -262,13 +262,13 @@ pub enum FfiMediaVerificationState { } #[cfg_attr(coverage_nightly, coverage(off))] -impl From<MediaVerificationState> for FfiMediaVerificationState { - fn from(value: MediaVerificationState) -> Self { +impl From<&Phase1InboundMediaState> for FfiMediaVerificationState { + fn from(value: &Phase1InboundMediaState) -> Self { match value { - MediaVerificationState::Pending => Self::Pending, - MediaVerificationState::Verified => Self::Verified, - MediaVerificationState::Failed => Self::Failed, - MediaVerificationState::Unavailable => Self::Unavailable, + Phase1InboundMediaState::Pending(_) => Self::Pending, + Phase1InboundMediaState::Verified(_) => Self::Verified, + Phase1InboundMediaState::Failed(_) => Self::Failed, + Phase1InboundMediaState::Unavailable => Self::Unavailable, } } } @@ -289,16 +289,17 @@ pub struct FfiMediaReferenceRecord { #[cfg_attr(coverage_nightly, coverage(off))] impl From<MediaReference> for FfiMediaReferenceRecord { fn from(value: MediaReference) -> Self { + let structural = value.structural(); Self { schema_version: MOBILE_FFI_SCHEMA_VERSION, - url: value.url, - sha256: value.sha256, - media_type: value.media_type, - width: value.width, - height: value.height, - byte_size: value.byte_size, - alt: value.alt, - verification: value.verification.into(), + url: structural.source_url().to_owned(), + sha256: structural.expected_sha256().map(str::to_owned), + media_type: structural.expected_media_type().map(str::to_owned), + width: structural.expected_width(), + height: structural.expected_height(), + byte_size: structural.expected_byte_size(), + alt: structural.alt().map(str::to_owned), + verification: value.retrieval().into(), } } } diff --git a/crates/mobile_ffi/src/error.rs b/crates/mobile_ffi/src/error.rs @@ -147,6 +147,7 @@ impl From<TodayError> for RadrootsAppError { ("today_cursor_invalid", true, &["restart_pagination"][..]) } TodayError::RuntimeUnavailable => ("today_runtime_unavailable", true, &["retry"][..]), + TodayError::InboundMedia(_) => ("today_media_invalid", false, &["retry_media"][..]), TodayError::CorruptProjection | TodayError::Serialization | TodayError::Storage(_) => { ("today_state_failed", true, &["rebuild", "retry"][..]) }