commit 28d6f6b901972bd524554b51531f23745c1d2b4f
parent 0ea3859ee11ee6b9e4181c9e11369e887905ae21
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:
6 files changed, 2225 insertions(+), 194 deletions(-)
diff --git a/core/crates/tera_core/Cargo.toml b/core/crates/tera_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/core/crates/tera_core/src/runtime/product_surface.rs b/core/crates/tera_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/core/crates/tera_core/src/runtime/product_surface/media.rs b/core/crates/tera_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/core/crates/tera_core/src/runtime/product_surface/model.rs b/core/crates/tera_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/core/crates/tera_core/src/runtime/product_surface/projection.rs b/core/crates/tera_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/core/crates/tera_core/src/runtime/product_surface/today.rs b/core/crates/tera_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, ¤t.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")