lib

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

commit 875c423d487cdc7c140bf7742c7986748ba7e602
parent dd429d5c16efe7c2f2ec1b79a6a316021d23023a
Author: triesap <tyson@radroots.org>
Date:   Sat,  8 Aug 2026 01:34:58 +0000

feat(mobile): sync Today from bounded relay data

- propagate exact selectors through canonical sync pulls
- fetch and typed-admit Today events before projection
- expose relay sync outcomes through the native boundary
- honor public, simulator, and device context policies

Diffstat:
Mcrates/mobile_core/src/runtime/product_surface.rs | 4+++-
Mcrates/mobile_core/src/runtime/product_surface/context.rs | 78++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++--
Mcrates/mobile_core/src/runtime/product_surface/today.rs | 226+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcrates/mobile_ffi/src/dto.rs | 75+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++--------------
Mcrates/mobile_ffi/src/runtime.rs | 73+++++++++++++++++++++++++++++++++++++++++++++++++++++++++----------------
Mcrates/mobile_ffi/tests/runtime_delegation.rs | 13+++++++++++++
Mcrates/sync/src/pull.rs | 23++++++++++++++++++++---
Mcrates/sync/tests/pull.rs | 39+++++++++++++++++++++++++++++++++++++++
8 files changed, 495 insertions(+), 36 deletions(-)

diff --git a/crates/mobile_core/src/runtime/product_surface.rs b/crates/mobile_core/src/runtime/product_surface.rs @@ -21,7 +21,7 @@ pub use authoring::{ }; pub use context::{ ContextAdmission, ContextRank, LocalNetwork, LocalNetworkAdmission, LocalNetworkError, - LocalityEvidence, + LocalNetworkRelayPolicy, LocalityEvidence, }; pub use cursor::{CursorError, CursorScope, TodayCursor, TodayCursorPosition}; pub use identity::{CARD_ID_SCHEMA_VERSION, CardId, CardIdError, CardSourceIdentity}; @@ -43,6 +43,8 @@ pub use ranking::{RankError, TODAY_RANK_SCHEMA_VERSION, TimeRelevance, TodayRank pub use today::{ TodayError, TodayIngestReceipt, TodayPageRequest, TodayProjectionUpdate, TodayRefreshReceipt, }; +#[cfg(feature = "mobile-social")] +pub use today::{TodayRelaySyncState, TodaySyncReceipt}; use super::RadrootsRuntime; diff --git a/crates/mobile_core/src/runtime/product_surface/context.rs b/crates/mobile_core/src/runtime/product_surface/context.rs @@ -35,6 +35,15 @@ pub enum LocalNetworkError { DuplicateAuthor, } +/// Host environment whose destination policy governs LocalNetwork relays. +#[derive(Clone, Copy, Debug, Deserialize, Eq, PartialEq, Serialize)] +#[serde(rename_all = "PascalCase")] +pub enum LocalNetworkRelayPolicy { + Public, + Simulator, + Device, +} + impl LocalNetwork { pub fn new( id: String, @@ -44,6 +53,28 @@ impl LocalNetwork { followed_authors: Vec<String>, generation: u64, ) -> Result<Self, LocalNetworkError> { + Self::new_for_relay_policy( + id, + label, + relay_urls, + locality, + followed_authors, + generation, + LocalNetworkRelayPolicy::Public, + ) + } + + /// Constructs a context under the exact host relay destination policy. + #[allow(clippy::too_many_arguments)] + pub fn new_for_relay_policy( + id: String, + label: String, + relay_urls: Vec<String>, + locality: Option<String>, + followed_authors: Vec<String>, + generation: u64, + relay_policy: LocalNetworkRelayPolicy, + ) -> Result<Self, LocalNetworkError> { validate_text(&id, "id")?; validate_text(&label, "label")?; if let Some(locality) = locality.as_deref() { @@ -58,8 +89,15 @@ impl LocalNetwork { if relay.is_empty() || relay.len() > RELAY_URL_MAX_BYTES { return Err(LocalNetworkError::InvalidRelay); } - let relay = RelayUrl::parse(relay, RelayUrlPolicy::Public) - .map_err(|_| LocalNetworkError::InvalidRelay)?; + let relay = RelayUrl::parse( + relay, + match relay_policy { + LocalNetworkRelayPolicy::Public => RelayUrlPolicy::Public, + LocalNetworkRelayPolicy::Simulator => RelayUrlPolicy::Local, + LocalNetworkRelayPolicy::Device => RelayUrlPolicy::PrivateNetwork, + }, + ) + .map_err(|_| LocalNetworkError::InvalidRelay)?; if !relays.insert(relay.clone()) { return Err(LocalNetworkError::DuplicateRelay); } @@ -369,5 +407,41 @@ mod tests { ), Err(LocalNetworkError::InvalidRelay) )); + assert!( + LocalNetwork::new_for_relay_policy( + "id".into(), + "label".into(), + vec!["ws://127.0.0.1:7447".into()], + None, + vec![], + 0, + LocalNetworkRelayPolicy::Simulator, + ) + .is_ok() + ); + assert!( + LocalNetwork::new_for_relay_policy( + "id".into(), + "label".into(), + vec!["wss://192.168.1.7:7447".into()], + None, + vec![], + 0, + LocalNetworkRelayPolicy::Device, + ) + .is_ok() + ); + assert!( + LocalNetwork::new_for_relay_policy( + "id".into(), + "label".into(), + vec!["ws://192.168.1.7:7447".into()], + None, + vec![], + 0, + LocalNetworkRelayPolicy::Device, + ) + .is_err() + ); } } diff --git a/crates/mobile_core/src/runtime/product_surface/today.rs b/crates/mobile_core/src/runtime/product_surface/today.rs @@ -19,6 +19,19 @@ use serde::{Deserialize, Serialize}; use sha2::{Digest, Sha256}; use thiserror::Error; +#[cfg(feature = "mobile-social")] +use radroots_event::admission::ContractValidatedEvent; +#[cfg(feature = "mobile-social")] +use radroots_sync::{ + PullRequest, + ingest::{AdmissionDecision, AdmissionPolicy}, + pull::PullTermination, +}; +#[cfg(feature = "mobile-social")] +use radroots_transport::{ + Target, outcome::FetchTargetState, source::FetchSelector, target::TargetSet, +}; + use super::{ CardId, CardLifecycleState, ClassifiedCard, CursorError, CursorScope, LocalAuthorOverlay, LocalNetwork, LocalityEvidence, MeSnapshot, MediaReference, MediaVerificationState, @@ -33,6 +46,12 @@ const TODAY_PROJECTION_DOCUMENT_SCHEMA_VERSION: u16 = 1; const TODAY_SNAPSHOT_SCHEMA_VERSION: u16 = 1; const TODAY_PAGE_LIMIT_MAX: u16 = 100; const TODAY_SEARCH_LIMIT_MAX: u16 = 100; +#[cfg(feature = "mobile-social")] +const TODAY_SYNC_PAGE_LIMIT: u16 = 500; +#[cfg(feature = "mobile-social")] +const TODAY_SYNC_MAX_PAGES: u16 = 8; +#[cfg(feature = "mobile-social")] +const TODAY_SYNC_KINDS: [u32; 7] = [0, 1, 5, 1111, 30_402, 31_922, 31_923]; const PROJECTION_GENERATION_DOMAIN: &[u8] = b"radroots.today-projection.v1\0"; const PROJECTION_CONTENT_DOMAIN: &[u8] = b"radroots.today-content-generation.v1\0"; const PROJECTION_DOCUMENT_KEY_DOMAIN: &[u8] = b"radroots.today-document-key.v1\0"; @@ -66,6 +85,27 @@ pub struct TodayIngestReceipt { pub projection: TodayRefreshReceipt, } +#[cfg(feature = "mobile-social")] +#[derive(Clone, Copy, Debug, Deserialize, Eq, PartialEq, Serialize)] +#[serde(rename_all = "PascalCase")] +pub enum TodayRelaySyncState { + Complete, + Partial, + Offline, +} + +#[cfg(feature = "mobile-social")] +#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)] +#[serde(rename_all = "camelCase")] +pub struct TodaySyncReceipt { + pub relay_state: TodayRelaySyncState, + pub pages_fetched: u16, + pub events_observed: u64, + pub events_admitted: u64, + pub events_rejected: u64, + pub projection: TodayRefreshReceipt, +} + #[derive(Clone, Debug, Eq, PartialEq)] pub struct TodayPageRequest { pub limit: u16, @@ -157,6 +197,72 @@ struct FrozenTodaySnapshot { } impl RadrootsRuntime { + /// Pulls bounded Today-relevant relay pages, canonically admits valid + /// observations, and materializes the selected LocalNetwork projection. + #[cfg(feature = "mobile-social")] + pub async fn phase1_sync_today( + &self, + context: &LocalNetwork, + now_unix_seconds: u64, + update: TodayProjectionUpdate, + ) -> Result<TodaySyncReceipt, TodayError> { + if now_unix_seconds == 0 { + return Err(TodayError::InvalidRequest); + } + let targets = context + .relay_urls + .iter() + .map(Target::nostr_relay) + .collect::<Result<Vec<_>, _>>() + .map_err(|_| TodayError::InvalidRequest)?; + let targets = TargetSet::new(targets).map_err(|_| TodayError::InvalidRequest)?; + let selector = FetchSelector::all() + .with_kinds(TODAY_SYNC_KINDS.to_vec()) + .map_err(|_| TodayError::InvalidRequest)?; + let request = PullRequest::new(targets, TODAY_SYNC_PAGE_LIMIT, TODAY_SYNC_MAX_PAGES) + .map_err(|_| TodayError::RuntimeUnavailable)? + .with_selector(selector); + let sync = self + .client + .sync() + .map_err(|_| TodayError::RuntimeUnavailable)? + .ok_or(TodayError::RuntimeUnavailable)?; + let pull = sync + .pull(request, &TodayAdmissionPolicy) + .await + .map_err(|_| TodayError::RuntimeUnavailable)?; + let projection = self + .phase1_refresh_today(context, now_unix_seconds, update) + .await?; + let events_admitted = pull + .ingest_outcomes() + .iter() + .filter(|outcome| outcome.is_ok()) + .count() as u64; + let events_observed = u64::try_from(pull.events_observed()).unwrap_or(u64::MAX); + let events_rejected = events_observed.saturating_sub(events_admitted); + let target_complete = !pull.target_outcomes().is_empty() + && pull + .target_outcomes() + .iter() + .all(|outcome| outcome.state() == FetchTargetState::Complete); + let relay_state = match pull.termination() { + PullTermination::Complete if target_complete => TodayRelaySyncState::Complete, + PullTermination::SourceFailed if pull.pages_fetched() == 0 => { + TodayRelaySyncState::Offline + } + _ => TodayRelaySyncState::Partial, + }; + Ok(TodaySyncReceipt { + relay_state, + pages_fetched: pull.pages_fetched(), + events_observed, + events_admitted, + events_rejected, + projection, + }) + } + /// Durably admits one already verified and visibility-authorized relay observation, /// then advances the selected LocalNetwork projection. pub async fn phase1_ingest_visible( @@ -544,6 +650,27 @@ impl RadrootsRuntime { } } +#[cfg(feature = "mobile-social")] +struct TodayAdmissionPolicy; + +#[cfg(feature = "mobile-social")] +impl AdmissionPolicy for TodayAdmissionPolicy { + fn policy_id(&self) -> &'static str { + "radroots.mobile.today.v1" + } + + fn decide(&self, event: &ContractValidatedEvent) -> AdmissionDecision { + let admitted = verify_nip01_event(event.event().clone()) + .ok() + .and_then(|event| admit_verified_event(event).ok()); + if admitted.is_some() { + AdmissionDecision::Visible + } else { + AdmissionDecision::Reject + } + } +} + fn ingest_receipt( receipt: AdmissionReceipt, projection: TodayRefreshReceipt, @@ -1121,6 +1248,9 @@ fn decode<T: for<'de> Deserialize<'de>>(value: &[u8]) -> Result<T, TodayError> { #[cfg(test)] mod tests { use super::*; + #[cfg(feature = "mobile-social")] + use std::sync::{Arc, Mutex, RwLock, atomic::AtomicBool}; + use nostr::secp256k1::Message; use nostr::{Keys, SECP256K1}; use radroots_event::{ @@ -1129,6 +1259,12 @@ mod tests { wire::{Nip01EventWire, compute_canonical_nip01_event_id}, }; use radroots_event_codec::verify::Nip01SignatureVerifier; + #[cfg(feature = "mobile-social")] + use radroots_transport::{ + Error as TransportError, EventSource, FetchPage, FetchRequest, SourceStatus, + outcome::{FetchTargetOutcome, FetchTargetState}, + source::NextPage, + }; use radroots_transport::{ Target, TransportId, source::{EventProvenance, ObservedEvent}, @@ -1138,6 +1274,48 @@ mod tests { struct Allow; + #[cfg(feature = "mobile-social")] + struct TodaySource { + event: SignedEvent, + requested_kinds: Mutex<Vec<Vec<u32>>>, + } + + #[cfg(feature = "mobile-social")] + impl EventSource for TodaySource { + fn status( + &self, + ) -> radroots_transport::BoxFuture<'_, Result<SourceStatus, TransportError>> { + Box::pin(async { unreachable!("Today sync does not inspect source status") }) + } + + fn fetch( + &self, + request: FetchRequest, + ) -> radroots_transport::BoxFuture<'_, Result<FetchPage, TransportError>> { + Box::pin(async move { + self.requested_kinds + .lock() + .expect("requested kinds") + .push(request.selector().kinds().to_vec()); + let target = request.target_set().targets()[0].clone(); + let provenance = EventProvenance::new( + TransportId::NOSTR, + target.fingerprint().clone(), + 2_000_000_100_000, + )?; + FetchPage::for_request( + &request, + vec![ObservedEvent::new(self.event.clone(), provenance)], + vec![FetchTargetOutcome::new( + target.fingerprint().clone(), + FetchTargetState::Complete, + )], + NextPage::Complete, + ) + }) + } + } + impl AdmissionPolicy for Allow { type Error = core::convert::Infallible; @@ -1292,6 +1470,54 @@ mod tests { .expect("ingest") } + #[cfg(feature = "mobile-social")] + #[tokio::test] + async fn relay_sync_fetches_the_exact_today_selector_and_projects_real_events() { + let source = Arc::new(TodaySource { + event: signed(1, Vec::new(), "Fresh from the field", 2_000_000_000), + requested_kinds: Mutex::new(Vec::new()), + }); + let client = radroots_sdk::ClientBuilder::memory_default() + .source(source.clone()) + .host_sync(radroots_sdk::sync::HostPolicy::standard()) + .build() + .expect("client"); + let runtime = RadrootsRuntime { + client, + started_unix_ms: 1, + shutting_down: AtomicBool::new(false), + platform_app: RwLock::new(None), + store_public_key: None, + }; + let context = context(None, 1); + + let receipt = runtime + .phase1_sync_today(&context, 2_000_000_200, TodayProjectionUpdate::Incremental) + .await + .expect("Today sync"); + assert_eq!(receipt.relay_state, TodayRelaySyncState::Complete); + assert_eq!(receipt.pages_fetched, 1); + assert_eq!(receipt.events_observed, 1); + assert_eq!(receipt.events_admitted, 1); + assert_eq!(receipt.events_rejected, 0); + assert_eq!(receipt.projection.visible_cards, 1); + assert_eq!( + source + .requested_kinds + .lock() + .expect("requested kinds") + .as_slice(), + &[TODAY_SYNC_KINDS.to_vec()] + ); + let page = runtime + .phase1_today_page(&context, TodayPageRequest::first(20, 2_000_000_200)) + .await + .expect("Today page"); + assert_eq!(page.items.len(), 1); + assert_eq!(page.items[0].card.card_type, TodayCardType::Update); + assert_eq!(page.items[0].card.content, "Fresh from the field"); + } + #[tokio::test] async fn equal_timestamp_pages_are_complete_and_remain_frozen_across_ingest() { let runtime = RadrootsRuntime::test_memory().expect("runtime"); diff --git a/crates/mobile_ffi/src/dto.rs b/crates/mobile_ffi/src/dto.rs @@ -19,12 +19,12 @@ use radroots_mobile_core::runtime::{ info::{AppInfo, RuntimeBuildInfo, RuntimeInfo}, product_surface::{ AddCommandType, CardLifecycleState, CreateAsk, CreateEvent, CreateFoodAvailability, - CreatePhotoUpdate, CreateUpdate, LocalNetwork, MeSnapshot, MediaReference, - MediaVerificationState, Phase1AddCommand, Phase1CancellationPolicy, Phase1DraftStatus, - Phase1MediaPrerequisite, Phase1MediaStage, Phase1OutboxState, Phase1QueuePolicy, - Phase1RelaySatisfaction, ProfileSummary, SearchResult, SearchResultType, SupportingProfile, - ThreadEntry, TodayCard, TodayCardType, TodayPage, TodayProjectionUpdate, - TodayRefreshReceipt, + CreatePhotoUpdate, CreateUpdate, LocalNetwork, LocalNetworkRelayPolicy, MeSnapshot, + MediaReference, MediaVerificationState, Phase1AddCommand, Phase1CancellationPolicy, + Phase1DraftStatus, Phase1MediaPrerequisite, Phase1MediaStage, Phase1OutboxState, + Phase1QueuePolicy, Phase1RelaySatisfaction, ProfileSummary, SearchResult, SearchResultType, + SupportingProfile, ThreadEntry, TodayCard, TodayCardType, TodayPage, TodayProjectionUpdate, + TodayRefreshReceipt, TodayRelaySyncState, TodaySyncReceipt, }, sdk::{ SdkCapabilityRecord, SdkRelayStatusRecord, SdkRelayStatusReportRecord, SdkShutdownRecord, @@ -210,14 +210,24 @@ impl TryFrom<FfiLocalNetworkRecord> for LocalNetwork { type Error = RadrootsAppError; fn try_from(value: FfiLocalNetworkRecord) -> Result<Self, Self::Error> { - require_schema(value.schema_version)?; - LocalNetwork::new( - value.id, - value.label, - value.relay_urls, - value.locality, - value.followed_authors, - value.generation, + value.try_into_with_relay_policy(LocalNetworkRelayPolicy::Public) + } +} + +impl FfiLocalNetworkRecord { + pub(crate) fn try_into_with_relay_policy( + self, + relay_policy: LocalNetworkRelayPolicy, + ) -> Result<LocalNetwork, RadrootsAppError> { + require_schema(self.schema_version)?; + LocalNetwork::new_for_relay_policy( + self.id, + self.label, + self.relay_urls, + self.locality, + self.followed_authors, + self.generation, + relay_policy, ) .map_err(|_| RadrootsAppError::invalid_argument("invalid_local_network")) } @@ -494,6 +504,43 @@ pub struct FfiTodayRefreshRecord { pub changed: bool, } +#[derive(Clone, Copy, Debug, Eq, PartialEq, uniffi::Enum)] +pub enum FfiTodayRelaySyncState { + Complete, + Partial, + Offline, +} + +#[derive(Clone, Debug, Eq, PartialEq, uniffi::Record)] +pub struct FfiTodaySyncRecord { + pub schema_version: u16, + pub relay_state: FfiTodayRelaySyncState, + pub pages_fetched: u16, + pub events_observed: u64, + pub events_admitted: u64, + pub events_rejected: u64, + pub projection: FfiTodayRefreshRecord, +} + +#[cfg_attr(coverage_nightly, coverage(off))] +impl From<TodaySyncReceipt> for FfiTodaySyncRecord { + fn from(value: TodaySyncReceipt) -> Self { + Self { + schema_version: MOBILE_FFI_SCHEMA_VERSION, + relay_state: match value.relay_state { + TodayRelaySyncState::Complete => FfiTodayRelaySyncState::Complete, + TodayRelaySyncState::Partial => FfiTodayRelaySyncState::Partial, + TodayRelaySyncState::Offline => FfiTodayRelaySyncState::Offline, + }, + pages_fetched: value.pages_fetched, + events_observed: value.events_observed, + events_admitted: value.events_admitted, + events_rejected: value.events_rejected, + projection: value.projection.into(), + } + } +} + #[cfg_attr(coverage_nightly, coverage(off))] impl From<TodayRefreshReceipt> for FfiTodayRefreshRecord { fn from(value: TodayRefreshReceipt) -> Self { diff --git a/crates/mobile_ffi/src/runtime.rs b/crates/mobile_ffi/src/runtime.rs @@ -1,5 +1,6 @@ use std::sync::Arc; +use radroots_mobile_core::runtime::product_surface::LocalNetworkRelayPolicy; use radroots_mobile_core::runtime::product_surface::TodayPageRequest; use crate::dto::PreparedMedia; @@ -11,7 +12,8 @@ use crate::{ FfiMeRecord, FfiQueuePolicyRecord, FfiRelayStatusReportRecord, FfiRuntimeChangeKind, FfiRuntimeInfoRecord, FfiSearchResultRecord, FfiShutdownRecord, FfiStorageStatusRecord, FfiSubscriptionHandle, FfiTodayPageRecord, FfiTodayProjectionUpdate, FfiTodayRefreshRecord, - RadrootsAppError, RadrootsHostSigner, RadrootsRuntimeObserver, add_schemas, decode_id, + FfiTodaySyncRecord, RadrootsAppError, RadrootsHostSigner, RadrootsRuntimeObserver, add_schemas, + decode_id, }; #[derive(Clone, Copy, Debug, Eq, PartialEq, uniffi::Enum)] @@ -253,7 +255,7 @@ impl RadrootsRuntime { &self, context: FfiLocalNetworkRecord, ) -> Result<FfiLocalNetworkRecord, RadrootsAppError> { - LocalNetworkRecordConversion::round_trip(context) + self.local_network(context).map(Into::into) } pub async fn phase1_today_page( @@ -268,7 +270,7 @@ impl RadrootsRuntime { "invalid_today_page_request", )); } - let context = context.try_into()?; + let context = self.local_network(context)?; let request = match cursor { Some(cursor) => TodayPageRequest::after(limit, cursor), None => TodayPageRequest::first( @@ -290,7 +292,7 @@ impl RadrootsRuntime { now_unix_s: u64, update: FfiTodayProjectionUpdate, ) -> Result<FfiTodayRefreshRecord, RadrootsAppError> { - let context = context.try_into()?; + let context = self.local_network(context)?; let receipt = self .inner .phase1_refresh_today(&context, now_unix_s, update.into()) @@ -300,6 +302,22 @@ impl RadrootsRuntime { Ok(receipt.into()) } + pub async fn phase1_sync_today( + &self, + context: FfiLocalNetworkRecord, + now_unix_s: u64, + update: FfiTodayProjectionUpdate, + ) -> Result<FfiTodaySyncRecord, RadrootsAppError> { + let context = self.local_network(context)?; + let receipt = self + .inner + .phase1_sync_today(&context, now_unix_s, update.into()) + .await + .map_err(RadrootsAppError::from)?; + self.subscriptions.notify(FfiRuntimeChangeKind::Today, None); + Ok(receipt.into()) + } + pub async fn phase1_search( &self, context: FfiLocalNetworkRecord, @@ -307,7 +325,7 @@ impl RadrootsRuntime { limit: u16, as_of_unix_s: u64, ) -> Result<Vec<FfiSearchResultRecord>, RadrootsAppError> { - let context = context.try_into()?; + let context = self.local_network(context)?; self.inner .phase1_search(&context, &query, limit, as_of_unix_s) .await @@ -320,7 +338,7 @@ impl RadrootsRuntime { context: FfiLocalNetworkRecord, as_of_unix_s: u64, ) -> Result<FfiMeRecord, RadrootsAppError> { - let context = context.try_into()?; + let context = self.local_network(context)?; let public_key = self .inner .authenticated_store_public_key_hex() @@ -514,6 +532,39 @@ impl RadrootsRuntime { } } +impl RadrootsRuntime { + fn local_network( + &self, + context: FfiLocalNetworkRecord, + ) -> Result<radroots_mobile_core::runtime::product_surface::LocalNetwork, RadrootsAppError> + { + let profile = self.inner.sdk_relay_status()?.ok_or_else(|| { + RadrootsAppError::failure( + "relay_profile_unavailable", + "relay", + true, + &["configure_relay"], + "The relay profile is unavailable.", + ) + })?; + let relay_policy = match profile.profile.as_str() { + "public" => LocalNetworkRelayPolicy::Public, + "simulator_local" => LocalNetworkRelayPolicy::Simulator, + "device_development" => LocalNetworkRelayPolicy::Device, + _ => { + return Err(RadrootsAppError::failure( + "relay_profile_unsupported", + "relay", + false, + &["configure_relay"], + "The relay profile is unsupported.", + )); + } + }; + context.try_into_with_relay_policy(relay_policy) + } +} + async fn build_runtime( application_support_directory: String, public_key_hex: String, @@ -544,13 +595,3 @@ async fn build_runtime( }) .map_err(Into::into) } - -struct LocalNetworkRecordConversion; - -impl LocalNetworkRecordConversion { - fn round_trip(value: FfiLocalNetworkRecord) -> Result<FfiLocalNetworkRecord, RadrootsAppError> { - let context: radroots_mobile_core::runtime::product_surface::LocalNetwork = - value.try_into()?; - Ok(context.into()) - } -} diff --git a/crates/mobile_ffi/tests/runtime_delegation.rs b/crates/mobile_ffi/tests/runtime_delegation.rs @@ -69,6 +69,19 @@ async fn native_boundary_delegates_the_complete_core_surface() { assert_eq!(simulator.relays[0].access, "read_write"); assert!( runtime + .phase1_local_network(FfiLocalNetworkRecord { + schema_version: MOBILE_FFI_SCHEMA_VERSION, + id: "simulator".to_owned(), + label: "Simulator".to_owned(), + relay_urls: vec!["ws://127.0.0.1:7447".to_owned()], + locality: None, + followed_authors: vec![], + generation: 1, + }) + .is_ok() + ); + assert!( + runtime .configure_simulator_relays(vec!["wss://relay.example".to_owned()]) .is_err() ); diff --git a/crates/sync/src/pull.rs b/crates/sync/src/pull.rs @@ -3,7 +3,7 @@ use radroots_transport::{ FetchRequest, outcome::FetchTargetOutcome, - source::{FETCH_PAGE_MAX_EVENTS, FetchBounds, FetchCursor, NextPage}, + source::{FETCH_PAGE_MAX_EVENTS, FetchBounds, FetchCursor, FetchSelector, NextPage}, target::TargetSet, }; @@ -24,6 +24,7 @@ pub struct PullRequest { page_limit: u16, max_pages: u16, cursor: Option<FetchCursor>, + selector: FetchSelector, } impl PullRequest { @@ -40,6 +41,7 @@ impl PullRequest { page_limit, max_pages, cursor: None, + selector: FetchSelector::all(), }) } @@ -64,6 +66,18 @@ impl PullRequest { pub const fn cursor(&self) -> Option<&FetchCursor> { self.cursor.as_ref() } + + /// Applies explicit transport-neutral event constraints to every page. + #[must_use] + pub fn with_selector(mut self, selector: FetchSelector) -> Self { + self.selector = selector; + self + } + + /// Returns the exact event constraints for this pull. + pub const fn selector(&self) -> &FetchSelector { + &self.selector + } } /// Deterministic reason why a bounded pull returned control to its caller. @@ -164,7 +178,8 @@ impl Engine { request.targets.clone(), bounds, ) - .map_err(|_| Error::InvalidPullRequest)?; + .map_err(|_| Error::InvalidPullRequest)? + .with_selector(request.selector.clone()); if let Some(current) = cursor.clone() { fetch = fetch.with_cursor(current); } @@ -250,6 +265,8 @@ impl<'de> serde::Deserialize<'de> for PullRequest { page_limit: u16, max_pages: u16, cursor: Option<FetchCursor>, + #[serde(default)] + selector: FetchSelector, } let wire = Wire::deserialize(deserializer)?; @@ -258,6 +275,6 @@ impl<'de> serde::Deserialize<'de> for PullRequest { if let Some(cursor) = wire.cursor { request = request.with_cursor(cursor); } - Ok(request) + Ok(request.with_selector(wire.selector)) } } diff --git a/crates/sync/tests/pull.rs b/crates/sync/tests/pull.rs @@ -38,6 +38,7 @@ struct RequestEvidence { cursor: Option<String>, limit: u16, deadline_unix_ms: u64, + kinds: Vec<u32>, } struct ScriptedSource { @@ -75,6 +76,7 @@ impl EventSource for ScriptedSource { cursor: request.cursor().map(|cursor| cursor.as_str().to_owned()), limit: request.bounds().limit(), deadline_unix_ms: request.bounds().deadline_unix_ms(), + kinds: request.selector().kinds().to_vec(), }); match self .responses @@ -200,6 +202,7 @@ fn single_and_multiple_pages_propagate_cursor_deadline_and_ingest_results() { assert_eq!(request.page_limit(), 20); assert_eq!(request.max_pages(), 1); assert!(request.cursor().is_none()); + assert!(request.selector().kinds().is_empty()); let receipt = block_on(single.pull(request, &RegistryPolicy::visible())).expect("pull"); assert_eq!(receipt.termination(), PullTermination::Complete); assert_eq!(receipt.pages_fetched(), 1); @@ -245,6 +248,42 @@ fn single_and_multiple_pages_propagate_cursor_deadline_and_ingest_results() { } #[test] +fn pull_propagates_the_exact_selector_to_every_page() { + let next = FetchCursor::parse("page-2").expect("cursor"); + let source = Arc::new(ScriptedSource::new(vec![ + Response::Page { + events: vec![], + state: FetchTargetState::Partial, + next: NextPage::Cursor(next), + }, + Response::Page { + events: vec![], + state: FetchTargetState::Complete, + next: NextPage::Complete, + }, + ])); + let pull = engine(source.clone(), Arc::new(FixedClock(100)), 50); + let selector = radroots_transport::source::FetchSelector::all() + .with_kinds(vec![0, 1, 5, 1111, 30402, 31922, 31923]) + .expect("selector"); + let request = PullRequest::new(targets(), 20, 2) + .expect("request") + .with_selector(selector.clone()); + assert_eq!(request.selector(), &selector); + + let receipt = block_on(pull.pull(request, &RegistryPolicy::visible())).expect("pull"); + assert_eq!(receipt.pages_fetched(), 2); + assert_eq!( + source + .requests() + .into_iter() + .map(|request| request.kinds) + .collect::<Vec<_>>(), + vec![selector.kinds().to_vec(), selector.kinds().to_vec()] + ); +} + +#[test] fn source_failure_and_cancelled_page_return_resumable_partial_receipts() { let cursor = FetchCursor::parse("resume").expect("cursor"); let source = Arc::new(ScriptedSource::new(vec![