lib

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

commit 18973e0aa62ae809b0495d5812b0c5ac8c7774ff
parent 828c462255b52a3d9bbfe8ef90b7bc38395bedfa
Author: triesap <tyson@radroots.org>
Date:   Mon, 10 Aug 2026 20:35:58 +0000

feat: authorize native mobile uploads

Diffstat:
Mcrates/mobile_core/src/runtime/builder.rs | 5+++--
Mcrates/mobile_core/src/runtime/product_surface.rs | 9+++++----
Mcrates/mobile_core/src/runtime/product_surface/outbox.rs | 196++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-
Mcrates/mobile_core/src/runtime/sdk.rs | 15++++++++++-----
Mcrates/mobile_ffi/src/dto.rs | 92+++++++++++++++++++++++++++++++++++++++++++++++++------------------------------
Mcrates/mobile_ffi/src/operations.rs | 7++-----
Mcrates/mobile_ffi/src/runtime.rs | 69+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcrates/mobile_ffi/tests/runtime_delegation.rs | 30+++++++++++++++++++++++-------
Mcrates/sdk/src/adapters/blossom.rs | 129+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcrates/sdk/src/transport.rs | 41+++++++++++++++++++++++++++++++++++++++++
10 files changed, 533 insertions(+), 60 deletions(-)

diff --git a/crates/mobile_core/src/runtime/builder.rs b/crates/mobile_core/src/runtime/builder.rs @@ -85,6 +85,7 @@ impl RuntimeBuilder { #[cfg(test)] mod tests { use super::RuntimeBuilder; + use crate::runtime::sdk::SdkRelayAccessRecord; use crate::runtime::store::{MobileUserStoreConfig, ProtectedDataAvailability}; const PUBLIC_KEY: &str = "79be667ef9dcbbac55a06295ce870b07029bfcdb2dce28d959f2815b16f81798"; @@ -128,7 +129,7 @@ mod tests { assert_eq!(report.state, "configured"); assert_eq!(report.relays.len(), 1); assert_eq!(report.relays[0].relay_url, "wss://radroots.org"); - assert_eq!(report.relays[0].access, "read_write"); + assert_eq!(report.relays[0].access, SdkRelayAccessRecord::ReadWrite); assert_eq!(report.relays[0].read_state, "unobserved"); assert_eq!(report.relays[0].write_state, "unobserved"); } @@ -154,7 +155,7 @@ mod tests { .expect("configured profile"); assert_eq!(report.profile, "simulator_local"); assert_eq!(report.relays.len(), 1); - assert_eq!(report.relays[0].access, "read_write"); + assert_eq!(report.relays[0].access, SdkRelayAccessRecord::ReadWrite); assert!( runtime .configure_public_relays(vec!["ws://127.0.0.1:8080".to_owned()]) diff --git a/crates/mobile_core/src/runtime/product_surface.rs b/crates/mobile_core/src/runtime/product_surface.rs @@ -47,10 +47,11 @@ pub use outbox::{ Phase1AddIntent, Phase1CancellationPolicy, Phase1DraftError, Phase1DraftEventTiming, Phase1DraftFormSnapshot, Phase1DraftKind, Phase1DraftMediaSnapshot, Phase1DraftStatus, Phase1ExistingDraft, Phase1MediaOrphanRecord, Phase1MediaPrerequisite, Phase1MediaStage, - Phase1OutboxState, Phase1ProfileStatus, Phase1QueueIntent, Phase1QueuePolicy, - Phase1RelaySatisfaction, Phase1ReviseIntent, Phase1RevisionPhase, Phase1RevisionPolicy, - Phase1RevisionStatus, Phase1RevisionTarget, Phase1UploadIntent, Phase1UploadPlan, - phase1_new_addressable_identifier, phase1_new_operation_id, phase1_operation_now_unix_ms, + Phase1NativeUploadJob, Phase1OutboxState, Phase1ProfileStatus, Phase1QueueIntent, + Phase1QueuePolicy, Phase1RelaySatisfaction, Phase1ReviseIntent, Phase1RevisionPhase, + Phase1RevisionPolicy, Phase1RevisionStatus, Phase1RevisionTarget, Phase1UploadIntent, + Phase1UploadPlan, phase1_new_addressable_identifier, phase1_new_operation_id, + phase1_operation_now_unix_ms, }; pub use projection::{ProductEventClassification, ProductEventExclusion, classify_admitted_event}; pub use ranking::{RankError, TODAY_RANK_SCHEMA_VERSION, TimeRelevance, TodayRank, TodayRankInput}; diff --git a/crates/mobile_core/src/runtime/product_surface/outbox.rs b/crates/mobile_core/src/runtime/product_surface/outbox.rs @@ -557,6 +557,39 @@ pub struct Phase1UploadPlan { pub updated_at_unix_ms: u64, } +/// Immutable Rust-authorized upload job handed to a native background +/// transfer implementation after the durable draft enters `media_uploading`. +#[derive(Clone)] +pub struct Phase1NativeUploadJob { + operation_id: [u8; 16], + remote_url: String, + authorization_header: String, + expected_sha256: String, + media_type: String, + byte_size: u64, +} + +impl Phase1NativeUploadJob { + pub const fn operation_id(&self) -> [u8; 16] { + self.operation_id + } + pub fn remote_url(&self) -> &str { + self.remote_url.as_str() + } + pub fn authorization_header(&self) -> &str { + self.authorization_header.as_str() + } + pub fn expected_sha256(&self) -> &str { + self.expected_sha256.as_str() + } + pub fn media_type(&self) -> &str { + self.media_type.as_str() + } + pub const fn byte_size(&self) -> u64 { + self.byte_size + } +} + /// Existing published card selected for one lossless revision operation. #[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)] #[serde(deny_unknown_fields)] @@ -590,6 +623,29 @@ impl Phase1RevisionTarget { Ok(value) } + /// Builds a revision target from its canonical source identity while + /// keeping Nostr kind parsing inside the Rust protocol boundary. + pub fn from_source( + command_type: AddCommandType, + card_id: CardId, + source_event_id: impl Into<String>, + source_address: Option<String>, + author_public_key: impl Into<String>, + ) -> Result<Self, Phase1DraftError> { + let source_kind = match source_address.as_deref() { + Some(address) => parse_address(address)?.0, + None => 1, + }; + Self::new( + command_type, + card_id, + source_event_id, + source_kind, + source_address, + author_public_key, + ) + } + fn validate(&self) -> Result<(), Phase1DraftError> { let author = PublicKey::from_hex(&self.author_public_key) .map_err(|_| Phase1DraftError::InvalidRevision)?; @@ -867,6 +923,7 @@ pub struct Phase1DraftStatus { state: Phase1OutboxState, card_id: CardId, push: Option<PushStatus>, + revision_policy: Option<Phase1RevisionPolicy>, } impl Phase1DraftStatus { @@ -894,6 +951,9 @@ impl Phase1DraftStatus { pub const fn push(&self) -> Option<&PushStatus> { self.push.as_ref() } + pub const fn revision_policy(&self) -> Option<Phase1RevisionPolicy> { + self.revision_policy + } } #[derive(Clone, Debug, Error, Eq, PartialEq)] @@ -1432,6 +1492,137 @@ impl RadrootsRuntime { .await } + /// Persists the upload transition and returns one immutable native job. + /// The authorization header is deliberately absent from durable draft + /// state and must never be persisted by the host. + pub async fn phase1_prepare_native_upload( + &self, + intent: Phase1UploadIntent, + ) -> Result<(Phase1DraftStatus, Phase1NativeUploadJob), Phase1DraftError> { + let now_unix_ms = phase1_operation_now_unix_ms()?; + let plan = Phase1UploadPlan::derive(now_unix_ms, phase1_random_id()?, phase1_random_id()?)?; + let request = radroots_sdk::transport::BlossomUploadRequest::new( + intent.bytes, + intent.media_type, + intent.dimensions, + now_unix_ms, + ) + .map_err(|_| Phase1DraftError::InvalidMedia)?; + let blossom = self + .client + .blossom() + .map_err(|_| Phase1DraftError::OperationUnavailable)? + .ok_or(Phase1DraftError::OperationUnavailable)?; + let transaction = blossom + .prepare_upload(request) + .map_err(|_| Phase1DraftError::Operation)?; + let remote_url = transaction.expected_url().as_str().to_owned(); + let content = radroots_blossom::authorization::AuthorizationContent::parse( + &plan.authorization_content, + ) + .map_err(|_| Phase1DraftError::InvalidMedia)?; + let claim = blossom + .authored_upload_claim( + &transaction, + content, + plan.authorization_created_at_unix_s, + plan.authorization_lifetime_seconds, + ) + .map_err(|_| Phase1DraftError::Operation)?; + let authorization = self + .phase1_authorize_blossom_upload( + plan.operation_id, + plan.artifact_id, + claim, + plan.signing_deadline_unix_ms, + plan.cancellation, + ) + .await?; + let uploading = self + .phase1_update_draft_media( + intent.draft_id, + intent.expected_revision, + remote_url.as_str(), + Phase1MediaStage::Uploading, + None, + now_unix_ms, + ) + .await?; + Ok(( + uploading, + Phase1NativeUploadJob { + operation_id: plan.operation_id, + remote_url, + authorization_header: authorization.into_string(), + expected_sha256: transaction.request().sha256().to_string(), + media_type: transaction.request().media_type().as_str().to_owned(), + byte_size: transaction.request().byte_size(), + }, + )) + } + + /// Verifies a native BUD-02 response and canonical BUD-01 retrieval before + /// advancing the durable media prerequisite. + pub async fn phase1_complete_native_upload( + &self, + intent: Phase1UploadIntent, + status_code: u16, + response_media_type: Option<&str>, + response_content_encoding: Option<&str>, + response_body: &[u8], + ) -> Result<Phase1DraftStatus, Phase1DraftError> { + let now_unix_ms = phase1_operation_now_unix_ms()?; + let request = radroots_sdk::transport::BlossomUploadRequest::new( + intent.bytes, + intent.media_type, + intent.dimensions, + now_unix_ms, + ) + .map_err(|_| Phase1DraftError::InvalidMedia)?; + let blossom = self + .client + .blossom() + .map_err(|_| Phase1DraftError::OperationUnavailable)? + .ok_or(Phase1DraftError::OperationUnavailable)?; + let transaction = blossom + .prepare_upload(request) + .map_err(|_| Phase1DraftError::Operation)?; + let url = transaction.expected_url().as_str().to_owned(); + match blossom + .complete_native_upload( + transaction, + status_code, + response_media_type, + response_content_encoding, + response_body, + radroots_sdk::transport::BlossomCancellation::default(), + ) + .await + { + Ok(receipt) => { + self.phase1_complete_draft_media( + intent.draft_id, + intent.expected_revision, + url.as_str(), + receipt, + now_unix_ms, + ) + .await + } + Err(error) => { + self.phase1_fail_draft_media( + intent.draft_id, + intent.expected_revision, + url.as_str(), + &error, + now_unix_ms, + ) + .await?; + Err(Phase1DraftError::Operation) + } + } + } + /// Cancels local work with a Rust-owned transition timestamp. pub async fn phase1_cancel_add_intent( &self, @@ -2653,6 +2844,7 @@ impl RadrootsRuntime { .target_card_id .map(Ok) .unwrap_or_else(|| card_id(payload.command_type, integrity.plan()))?; + let revision_policy = payload.revision.as_ref().map(|value| value.policy); let push = self.push_status_for(&draft).await?; if draft.stage() == AuthoredDraftStage::Queued && push.is_none() { return Err(Phase1DraftError::Corrupt); @@ -2667,6 +2859,7 @@ impl RadrootsRuntime { state, card_id, push, + revision_policy, }) } @@ -3377,11 +3570,10 @@ mod tests { let identifier = "market:summer:2026"; let source = CardSourceIdentity::address(31_923, AUTHOR, identifier).unwrap(); let target_card = CardId::derive(TodayCardType::Event, &source); - let target = Phase1RevisionTarget::new( + let target = Phase1RevisionTarget::from_source( AddCommandType::CreateEvent, target_card, "b".repeat(64), - 31_923, Some(format!("31923:{AUTHOR}:{identifier}")), AUTHOR, ) diff --git a/crates/mobile_core/src/runtime/sdk.rs b/crates/mobile_core/src/runtime/sdk.rs @@ -21,9 +21,15 @@ pub struct SdkStorageStatusRecord { } #[derive(Clone, Debug, Eq, PartialEq)] +pub enum SdkRelayAccessRecord { + ReadOnly, + ReadWrite, +} + +#[derive(Clone, Debug, Eq, PartialEq)] pub struct SdkRelayStatusRecord { pub relay_url: String, - pub access: String, + pub access: SdkRelayAccessRecord, pub read_state: String, pub write_state: String, pub read_last_attempt_unix_ms: Option<u64>, @@ -264,11 +270,10 @@ impl RadrootsRuntime { .map(|relay| SdkRelayStatusRecord { relay_url: relay.endpoint().url().to_string(), access: if relay.endpoint().access().can_write() { - "read_write" + SdkRelayAccessRecord::ReadWrite } else { - "read_only" - } - .to_owned(), + SdkRelayAccessRecord::ReadOnly + }, read_state: relay_evidence_label(relay.read().state()).to_owned(), write_state: relay_evidence_label(relay.write().state()).to_owned(), read_last_attempt_unix_ms: relay.read().last_attempt_unix_ms(), diff --git a/crates/mobile_ffi/src/dto.rs b/crates/mobile_ffi/src/dto.rs @@ -32,7 +32,7 @@ use radroots_mobile_core::runtime::{ }, sdk::{ SdkBlossomConfigurationRecord, SdkBlossomEvidenceRecord, SdkCapabilityRecord, - SdkRelayStatusRecord, SdkRelayStatusReportRecord, SdkShutdownRecord, + SdkRelayAccessRecord, SdkRelayStatusRecord, SdkRelayStatusReportRecord, SdkShutdownRecord, SdkStorageStatusRecord, }, }; @@ -658,6 +658,7 @@ pub struct FfiAddFieldRecord { pub required: bool, pub choices: Vec<String>, pub max_bytes: Option<u64>, + pub max_items: Option<u16>, } #[derive(Clone, Debug, Eq, PartialEq, uniffi::Record)] @@ -681,7 +682,19 @@ pub fn add_schemas() -> Vec<FfiAddSchemaRecord> { required, choices: choices.iter().map(|value| (*value).to_owned()).collect(), max_bytes, + max_items: None, }; + let media_field = |label: &str, required: bool, max_items| FfiAddFieldRecord { + max_items: Some(max_items), + ..field( + "media", + label, + Kind::Media, + required, + &[], + Some(MEDIA_FILE_MAX_BYTES), + ) + }; vec![ FfiAddSchemaRecord { schema_version: MOBILE_FFI_SCHEMA_VERSION, @@ -709,14 +722,7 @@ pub fn add_schemas() -> Vec<FfiAddSchemaRecord> { &[], Some(65_535), ), - field( - "media", - "Photos", - Kind::Media, - true, - &[], - Some(MEDIA_FILE_MAX_BYTES), - ), + media_field("Photos", true, 20), ], }, FfiAddSchemaRecord { @@ -732,14 +738,7 @@ pub fn add_schemas() -> Vec<FfiAddSchemaRecord> { &[], Some(65_535), ), - field( - "media", - "Photos", - Kind::Media, - false, - &[], - Some(MEDIA_FILE_MAX_BYTES), - ), + media_field("Photos", false, 20), ], }, FfiAddSchemaRecord { @@ -767,14 +766,7 @@ pub fn add_schemas() -> Vec<FfiAddSchemaRecord> { &[], Some(256), ), - field( - "media", - "Photo", - Kind::Media, - false, - &[], - Some(MEDIA_FILE_MAX_BYTES), - ), + media_field("Photo", false, 1), ], }, FfiAddSchemaRecord { @@ -821,14 +813,7 @@ pub fn add_schemas() -> Vec<FfiAddSchemaRecord> { &[], Some(64), ), - field( - "media", - "Photos", - Kind::Media, - false, - &[], - Some(MEDIA_FILE_MAX_BYTES), - ), + media_field("Photos", false, 20), ], }, ] @@ -915,6 +900,32 @@ pub struct FfiBlossomUploadIntent { pub media: FfiPreparedMediaInput, } +/// Secret-bearing native job. Hosts may use the authorization header for the +/// immediate OS request but must not persist it in application metadata. +#[derive(Clone, Debug, Eq, PartialEq, uniffi::Record)] +pub struct FfiNativeUploadJobRecord { + pub schema_version: u16, + pub operation_id: String, + pub draft: FfiDraftStatusRecord, + pub remote_url: String, + pub authorization_header: String, + pub expected_sha256: String, + pub media_type: String, + pub byte_size: u64, +} + +#[derive(Clone, Debug, Eq, PartialEq, uniffi::Record)] +pub struct FfiNativeUploadCompletionInput { + pub schema_version: u16, + pub draft_id: String, + pub expected_revision: u64, + pub media: FfiPreparedMediaInput, + pub status_code: u16, + pub response_media_type: Option<String>, + pub response_content_encoding: Option<String>, + pub response_body: Vec<u8>, +} + impl FfiAddDraftInput { pub(crate) fn command_and_media( self, @@ -1647,6 +1658,7 @@ pub struct FfiDraftStatusRecord { pub updated_at_unix_ms: u64, pub media: Vec<FfiDraftMediaRecord>, pub settlement: Option<FfiOperationSettlementRecord>, + pub is_revision: bool, } #[cfg_attr(coverage_nightly, coverage(off))] @@ -1669,6 +1681,7 @@ impl From<Phase1DraftStatus> for FfiDraftStatusRecord { updated_at_unix_ms: draft.updated_at_unix_ms(), media: value.media().iter().map(Into::into).collect(), settlement: settlement.map(Into::into), + is_revision: value.revision_policy().is_some(), } } } @@ -1822,11 +1835,17 @@ impl From<SdkStorageStatusRecord> for FfiStorageStatusRecord { } } +#[derive(Clone, Copy, Debug, Eq, PartialEq, uniffi::Enum)] +pub enum FfiRelayAccessRecord { + ReadOnly, + ReadWrite, +} + #[derive(Clone, Debug, Eq, PartialEq, uniffi::Record)] pub struct FfiRelayStatusRecord { pub schema_version: u16, pub relay_url: String, - pub access: String, + pub access: FfiRelayAccessRecord, pub read_state: String, pub write_state: String, pub read_last_attempt_unix_ms: Option<u64>, @@ -1841,7 +1860,10 @@ impl From<SdkRelayStatusRecord> for FfiRelayStatusRecord { Self { schema_version: MOBILE_FFI_SCHEMA_VERSION, relay_url: value.relay_url, - access: value.access, + access: match value.access { + SdkRelayAccessRecord::ReadOnly => FfiRelayAccessRecord::ReadOnly, + SdkRelayAccessRecord::ReadWrite => FfiRelayAccessRecord::ReadWrite, + }, read_state: value.read_state, write_state: value.write_state, read_last_attempt_unix_ms: value.read_last_attempt_unix_ms, diff --git a/crates/mobile_ffi/src/operations.rs b/crates/mobile_ffi/src/operations.rs @@ -491,10 +491,8 @@ impl From<Phase1ProfileStatus> for FfiProfileStatusRecord { #[derive(Clone, Debug, uniffi::Record)] pub struct FfiRevisionInputRecord { pub schema_version: u16, - pub command_type: crate::FfiAddCommandType, pub card_id: String, pub source_event_id: String, - pub source_kind: u32, pub source_address: Option<String>, pub author_public_key: String, pub replacement: FfiAddDraftInput, @@ -503,12 +501,11 @@ pub struct FfiRevisionInputRecord { impl FfiRevisionInputRecord { pub(crate) fn target(&self) -> Result<Phase1RevisionTarget, RadrootsAppError> { require_schema(self.schema_version)?; - Phase1RevisionTarget::new( - self.command_type.into(), + Phase1RevisionTarget::from_source( + self.replacement.command_type.into(), radroots_mobile_core::runtime::product_surface::CardId::parse(&self.card_id) .map_err(|_| RadrootsAppError::invalid_argument("invalid_card_id"))?, self.source_event_id.clone(), - self.source_kind, self.source_address.clone(), self.author_public_key.clone(), ) diff --git a/crates/mobile_ffi/src/runtime.rs b/crates/mobile_ffi/src/runtime.rs @@ -733,6 +733,75 @@ impl RadrootsRuntime { Ok(status.into()) } + /// Persists the upload transition before returning an immutable native + /// background-transfer job. + pub async fn phase1_prepare_add_media_background( + &self, + input: crate::FfiBlossomUploadIntent, + ) -> Result<crate::FfiNativeUploadJobRecord, RadrootsAppError> { + if input.schema_version != crate::MOBILE_FFI_SCHEMA_VERSION { + return Err(RadrootsAppError::invalid_argument( + "unsupported_schema_version", + )); + } + let draft_id = decode_id(&input.draft_id, "invalid_draft_id")?; + let intent = PreparedMedia::try_from(input.media)? + .into_upload_intent(draft_id, input.expected_revision)?; + let (status, job) = self + .inner + .phase1_prepare_native_upload(intent) + .await + .map_err(RadrootsAppError::from)?; + self.subscriptions + .notify(FfiRuntimeChangeKind::Media, Some(input.draft_id.clone())); + self.subscriptions + .notify(FfiRuntimeChangeKind::Drafts, Some(input.draft_id)); + Ok(crate::FfiNativeUploadJobRecord { + schema_version: crate::MOBILE_FFI_SCHEMA_VERSION, + operation_id: hex::encode(job.operation_id()), + draft: status.into(), + remote_url: job.remote_url().to_owned(), + authorization_header: job.authorization_header().to_owned(), + expected_sha256: job.expected_sha256().to_owned(), + media_type: job.media_type().to_owned(), + byte_size: job.byte_size(), + }) + } + + /// Accepts only bounded native HTTP evidence; Rust performs descriptor and + /// exact-byte retrieval verification before advancing durable state. + pub async fn phase1_complete_add_media_background( + &self, + input: crate::FfiNativeUploadCompletionInput, + ) -> Result<FfiDraftStatusRecord, RadrootsAppError> { + if input.schema_version != crate::MOBILE_FFI_SCHEMA_VERSION + || input.response_body.len() > 16_384 + { + return Err(RadrootsAppError::invalid_argument( + "invalid_native_upload_completion", + )); + } + let draft_id = decode_id(&input.draft_id, "invalid_draft_id")?; + let intent = PreparedMedia::try_from(input.media)? + .into_upload_intent(draft_id, input.expected_revision)?; + let status = self + .inner + .phase1_complete_native_upload( + intent, + input.status_code, + input.response_media_type.as_deref(), + input.response_content_encoding.as_deref(), + input.response_body.as_slice(), + ) + .await + .map_err(RadrootsAppError::from)?; + self.subscriptions + .notify(FfiRuntimeChangeKind::Media, Some(input.draft_id.clone())); + self.subscriptions + .notify(FfiRuntimeChangeKind::Drafts, Some(input.draft_id)); + Ok(status.into()) + } + pub async fn phase1_cancel_draft( &self, draft_id: String, diff --git a/crates/mobile_ffi/tests/runtime_delegation.rs b/crates/mobile_ffi/tests/runtime_delegation.rs @@ -3,8 +3,9 @@ use radroots_mobile_ffi::{ FfiBlossomUploadIntent, FfiCancellationPolicy, FfiDraftKind, FfiIdentityCommandKind, FfiIdentityCommandRecord, FfiIdentityLockState, FfiLocalNetworkRecord, FfiOutboxState, FfiPreparedMediaInput, FfiProfileMetadataInputRecord, FfiQueuePolicyRecord, - FfiRelaySatisfaction, FfiRetractionDraftInput, FfiRevisionInputRecord, FfiRevisionPhase, - FfiTodayCardType, FfiTodayProjectionUpdate, MOBILE_FFI_SCHEMA_VERSION, RadrootsAppError, + FfiRelayAccessRecord, FfiRelaySatisfaction, FfiRetractionDraftInput, FfiRevisionInputRecord, + FfiRevisionPhase, FfiTodayCardType, FfiTodayProjectionUpdate, MOBILE_FFI_SCHEMA_VERSION, + RadrootsAppError, }; mod support; @@ -40,7 +41,7 @@ async fn native_boundary_delegates_the_complete_core_surface() { assert_eq!(public.read_availability, "unavailable"); assert_eq!(public.write_availability, "unavailable"); assert_eq!(public.relays.len(), 1); - assert_eq!(public.relays[0].access, "read_write"); + assert_eq!(public.relays[0].access, FfiRelayAccessRecord::ReadWrite); assert_eq!(public.relays[0].read_state, "unobserved"); assert_eq!(public.relays[0].write_state, "unobserved"); @@ -52,7 +53,7 @@ async fn native_boundary_delegates_the_complete_core_surface() { .expect("relay status") .expect("public profile"); assert_eq!(public.relays.len(), 2); - assert_eq!(public.relays[1].access, "read_write"); + assert_eq!(public.relays[1].access, FfiRelayAccessRecord::ReadWrite); assert!( runtime .configure_public_relays(vec!["ws://127.0.0.1:7447".to_owned()]) @@ -68,7 +69,7 @@ async fn native_boundary_delegates_the_complete_core_surface() { .expect("simulator profile"); assert_eq!(simulator.profile, "simulator_local"); assert_eq!(simulator.relays.len(), 1); - assert_eq!(simulator.relays[0].access, "read_write"); + assert_eq!(simulator.relays[0].access, FfiRelayAccessRecord::ReadWrite); assert!( runtime .phase1_local_network(FfiLocalNetworkRecord { @@ -183,6 +184,22 @@ async fn native_boundary_delegates_the_complete_core_surface() { for (index, item) in parity.iter().enumerate() { assert_eq!(item.command_type, schemas[index].command_type); } + assert_eq!( + schemas[1] + .fields + .iter() + .find(|field| field.id == "media") + .and_then(|field| field.max_items), + Some(20) + ); + assert_eq!( + schemas[3] + .fields + .iter() + .find(|field| field.id == "media") + .and_then(|field| field.max_items), + Some(1) + ); let local_network = runtime .phase1_local_network(FfiLocalNetworkRecord { schema_version: MOBILE_FFI_SCHEMA_VERSION, @@ -484,10 +501,8 @@ async fn native_boundary_delegates_the_complete_core_surface() { let revision = runtime .phase1_save_revision_intent(FfiRevisionInputRecord { schema_version: MOBILE_FFI_SCHEMA_VERSION, - command_type: FfiAddCommandType::CreateUpdate, card_id, source_event_id, - source_kind: 1, source_address: None, author_public_key: support::PUBLIC_KEY.to_owned(), replacement: FfiAddDraftInput { @@ -517,6 +532,7 @@ async fn native_boundary_delegates_the_complete_core_surface() { .expect("save lossless revision intent"); assert_eq!(revision.phase, FfiRevisionPhase::ReplacementPending); assert_eq!(revision.operation_id, revision.replacement.draft_id); + assert!(revision.replacement.is_revision); let cancelled_revision = runtime .phase1_cancel_revision(revision.operation_id.clone()) .await diff --git a/crates/sdk/src/adapters/blossom.rs b/crates/sdk/src/adapters/blossom.rs @@ -336,6 +336,85 @@ async fn upload_bound_with_authorization( )) } +pub(crate) async fn complete_native_upload( + transaction: BlossomUploadTransaction, + status_code: u16, + response_media_type: Option<&str>, + response_content_encoding: Option<&str>, + response_body: &[u8], + cancellation: BlossomCancellation, +) -> Result<BlossomUploadReceipt, BlossomError> { + let config = transaction.config().clone(); + let expected_url = transaction.expected_url().clone(); + let request = transaction.into_request(); + if !matches!(status_code, 200 | 201) { + let status = StatusCode::from_u16(status_code).unwrap_or(StatusCode::INTERNAL_SERVER_ERROR); + return Err(http_status_error(status, BlossomPhase::Upload, true, 1)); + } + if response_media_type != Some("application/json") + || response_content_encoding.is_some_and(|value| value != "identity") + || response_body.is_empty() + || response_body.len() > config.max_descriptor_bytes() + { + return Err(failure( + BlossomErrorKind::InvalidDescriptor, + BlossomPhase::Upload, + false, + true, + 1, + )); + } + let descriptor = serde_json::from_slice::<BlobDescriptor>(response_body).map_err(|_| { + failure( + BlossomErrorKind::InvalidDescriptor, + BlossomPhase::Upload, + false, + true, + 1, + ) + })?; + let verified_upload = verify_descriptor(&request, &expected_url, descriptor, 1)?; + let retrieved = retrieve( + config, + BlossomInboundRequest::new( + verified_upload.url().as_blob_url().clone(), + Some(request.media_type().clone()), + Some(request.byte_size()), + Some(request.dimensions()), + )?, + cancellation, + ) + .await?; + if retrieved.bytes() != request.bytes() { + return Err(failure( + BlossomErrorKind::RetrievedBytesMismatch, + BlossomPhase::Verification, + false, + true, + 1_u8.saturating_add(retrieved.attempts()), + )); + } + let descriptor = verified_upload + .into_descriptor() + .approve_reference() + .and_then(|approved| approved.verify_bytes(retrieved.bytes(), request.media_type())) + .map_err(|_| { + failure( + BlossomErrorKind::RetrievedBytesMismatch, + BlossomPhase::Verification, + false, + true, + 1_u8.saturating_add(retrieved.attempts()), + ) + })?; + Ok(BlossomUploadReceipt::new( + descriptor, + request.dimensions(), + 1_u8.saturating_add(retrieved.attempts()), + request.verified_at_unix_ms(), + )) +} + pub(crate) async fn retrieve( config: BlossomConfig, request: BlossomInboundRequest, @@ -2029,6 +2108,56 @@ mod tests { .expect("profile") } + #[tokio::test] + async fn native_upload_completion_verifies_descriptor_and_retrieved_bytes() { + let bytes = png(2, 3); + let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); + let origin = format!("http://{}", listener.local_addr().unwrap()); + let request = upload_request(origin.as_str(), bytes.clone()); + let expected_url = + BlobUrl::parse(format!("{origin}/{}.png", request.sha256()).as_str()).unwrap(); + let descriptor = BlobDescriptor::new( + expected_url, + request.sha256(), + request.byte_size(), + MediaType::parse("image/png").unwrap(), + 1_900_000_000, + ) + .unwrap(); + let descriptor_body = serde_json::to_vec(&descriptor).unwrap(); + let server = tokio::spawn(async move { + let (mut stream, _) = listener.accept().await.unwrap(); + let request = read_request(&mut stream).await; + assert!(String::from_utf8_lossy(&request).starts_with("GET /")); + let head = format!( + "HTTP/1.1 200 OK\r\nContent-Type: image/png\r\nContent-Length: {}\r\nConnection: close\r\n\r\n", + bytes.len() + ); + stream.write_all(head.as_bytes()).await.unwrap(); + stream.write_all(&bytes).await.unwrap(); + stream.shutdown().await.unwrap(); + }); + let slot = crate::transport::BlossomSlot::new(); + slot.configure(config(origin.as_str())).unwrap(); + let transaction = slot + .prepare_upload(upload_request(origin.as_str(), png(2, 3))) + .unwrap(); + let receipt = slot + .complete_native_upload( + transaction, + 200, + Some("application/json"), + None, + descriptor_body.as_slice(), + BlossomCancellation::default(), + ) + .await + .unwrap(); + assert_eq!(receipt.descriptor().sha256(), descriptor.sha256()); + assert_eq!(receipt.attempts(), 2); + server.await.unwrap(); + } + async fn spawn_inbound_server( bytes: Vec<u8>, content_type: &'static str, diff --git a/crates/sdk/src/transport.rs b/crates/sdk/src/transport.rs @@ -1453,6 +1453,47 @@ impl BlossomSlot { } } + /// Verifies a host-executed BUD-02 response, then performs the canonical + /// BUD-01 exact-byte retrieval before returning an upload receipt. + pub async fn complete_native_upload( + &self, + transaction: BlossomUploadTransaction, + status_code: u16, + response_media_type: Option<&str>, + response_content_encoding: Option<&str>, + response_body: &[u8], + cancellation: BlossomCancellation, + ) -> Result<BlossomUploadReceipt, BlossomError> { + self.validate_transaction(&transaction)?; + let fingerprint = transaction.config_fingerprint(); + let result = crate::adapters::blossom::complete_native_upload( + transaction, + status_code, + response_media_type, + response_content_encoding, + response_body, + cancellation, + ) + .await; + match result { + Ok(receipt) => { + self.record_evidence(fingerprint, BlossomPhase::Verification, true, |evidence| { + evidence.record_success(BlossomEvidenceState::RetrievalVerified, None); + })?; + Ok(receipt) + } + Err(error) => { + self.record_evidence(fingerprint, BlossomPhase::Verification, true, |evidence| { + if error.possible_orphan() { + evidence.last_successful_state = BlossomEvidenceState::UploadVerified; + } + evidence.record_failure(&error); + })?; + Err(error) + } + } + } + /// Retrieves and verifies one immutable BUD-01 image under the exact /// configured DNS, TLS, redirect, retry, and byte limits. pub async fn retrieve(