commit 77aa55c0894c0142dbae43e272a515fd093eac36 parent 858b9e0d4fc25c3bc53ca9550c898fe8148c1bd8 Author: triesap <tyson@radroots.org> Date: Mon, 21 Sep 2026 14:20:09 +0000 transport: constrain delivery attempts to selected targets - Preserve frozen full request and receipt bindings - Hold unselected targets without transport admission - Retain shared claims and late acceptance across stop and reload - Verify public APIs, owner coverage and workspace consumers Diffstat:
22 files changed, 961 insertions(+), 22 deletions(-)
diff --git a/contracts/api_baselines/radroots_sdk.txt b/contracts/api_baselines/radroots_sdk.txt @@ -362,6 +362,7 @@ impl<'a> radroots_sdk::sync::Operations<'a> pub async fn radroots_sdk::sync::Operations<'a>::admit_signed(&self, radroots_sync::policy::SyncId) -> core::result::Result<radroots_sync::push::AdmissionRunReceipt, radroots_sync::policy::Error> pub async fn radroots_sdk::sync::Operations<'a>::cancel_push(&self, radroots_sync::policy::SyncId) -> core::result::Result<radroots_sync::push::PushCancellationReceipt, radroots_sync::policy::Error> pub async fn radroots_sdk::sync::Operations<'a>::deliver_push(&self, radroots_sync::policy::SyncId) -> core::result::Result<radroots_sync::push::DeliveryExecutionReceipt, radroots_sync::policy::Error> +pub async fn radroots_sdk::sync::Operations<'a>::deliver_push_selected(&self, radroots_sync::policy::SyncId, radroots_transport::target::TargetSet) -> core::result::Result<radroots_sync::push::DeliveryExecutionReceipt, radroots_sync::policy::Error> pub async fn radroots_sdk::sync::Operations<'a>::ingest(&self, radroots_transport::source::ObservedEvent, &dyn radroots_sync::ingest::AdmissionPolicy) -> core::result::Result<radroots_sync::ingest::IngestReceipt, radroots_sync::policy::Error> pub async fn radroots_sdk::sync::Operations<'a>::ingest_batch(&self, alloc::vec::Vec<radroots_transport::source::ObservedEvent>, &dyn radroots_sync::ingest::AdmissionPolicy) -> radroots_sync::ingest::IngestBatchReceipt pub async fn radroots_sdk::sync::Operations<'a>::prepare_push(&self, radroots_sync::push::PushRequest) -> core::result::Result<radroots_sync::push::PushPreparation, radroots_sync::policy::Error> @@ -695,6 +696,7 @@ impl core::fmt::Debug for radroots_sdk::transport::NostrSlot pub fn radroots_sdk::transport::NostrSlot::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result impl radroots_transport::sink::EventSink for radroots_sdk::transport::NostrSlot pub fn radroots_sdk::transport::NostrSlot::deliver(&self, radroots_transport::sink::DeliveryRequest) -> radroots_transport::source::BoxFuture<'_, core::result::Result<radroots_transport::sink::DeliveryReceipt, radroots_transport::sink::SinkFailure>> +pub fn radroots_sdk::transport::NostrSlot::deliver_selected(&self, radroots_transport::sink::DeliveryRequest, radroots_transport::target::TargetSet) -> radroots_transport::source::BoxFuture<'_, core::result::Result<radroots_transport::sink::DeliveryReceipt, radroots_transport::sink::SinkFailure>> pub fn radroots_sdk::transport::NostrSlot::status(&self) -> radroots_transport::source::BoxFuture<'_, core::result::Result<radroots_transport::status::SinkStatus, radroots_transport::error::Error>> impl radroots_transport::source::EventSource for radroots_sdk::transport::NostrSlot pub fn radroots_sdk::transport::NostrSlot::fetch(&self, radroots_transport::source::FetchRequest) -> radroots_transport::source::BoxFuture<'_, core::result::Result<radroots_transport::source::FetchPage, radroots_transport::error::Error>> diff --git a/contracts/api_baselines/radroots_sync.txt b/contracts/api_baselines/radroots_sync.txt @@ -304,6 +304,7 @@ pub fn radroots_sync::Engine::source(&self) -> core::option::Option<&dyn radroot pub fn radroots_sync::Engine::storage(&self) -> &dyn radroots_sync::policy::SyncStorage impl radroots_sync::Engine pub async fn radroots_sync::Engine::deliver_push(&self, radroots_sync::policy::SyncId) -> core::result::Result<radroots_sync::push::DeliveryExecutionReceipt, radroots_sync::policy::Error> +pub async fn radroots_sync::Engine::deliver_push_selected(&self, radroots_sync::policy::SyncId, radroots_transport::target::TargetSet) -> core::result::Result<radroots_sync::push::DeliveryExecutionReceipt, radroots_sync::policy::Error> impl radroots_sync::Engine pub async fn radroots_sync::Engine::ingest(&self, radroots_transport::source::ObservedEvent, &dyn radroots_sync::ingest::AdmissionPolicy) -> core::result::Result<radroots_sync::ingest::IngestReceipt, radroots_sync::policy::Error> pub async fn radroots_sync::Engine::ingest_batch(&self, alloc::vec::Vec<radroots_transport::source::ObservedEvent>, &dyn radroots_sync::ingest::AdmissionPolicy) -> radroots_sync::ingest::IngestBatchReceipt diff --git a/contracts/api_baselines/radroots_transport.txt b/contracts/api_baselines/radroots_transport.txt @@ -121,6 +121,7 @@ pub radroots_transport::error::Error::FetchSelectorTooLarge pub radroots_transport::error::Error::InvalidDeliveryDeadline pub radroots_transport::error::Error::InvalidDeliveryOutcome pub radroots_transport::error::Error::InvalidDeliveryRequestId +pub radroots_transport::error::Error::InvalidDeliveryTargetSelection pub radroots_transport::error::Error::InvalidDeliveryTimestamp pub radroots_transport::error::Error::InvalidFetchCursor pub radroots_transport::error::Error::InvalidFetchDeadline @@ -260,6 +261,7 @@ pub const fn radroots_transport::sink::DeliveryRequest::payload(&self) -> &radro pub const fn radroots_transport::sink::DeliveryRequest::request_id(&self) -> &radroots_transport::sink::DeliveryRequestId pub const fn radroots_transport::sink::DeliveryRequest::satisfaction(&self) -> &radroots_transport::policy::SatisfactionPolicy pub const fn radroots_transport::sink::DeliveryRequest::target_set(&self) -> &radroots_transport::target::TargetSet +pub fn radroots_transport::sink::DeliveryRequest::validate_target_selection(&self, &radroots_transport::target::TargetSet) -> core::result::Result<(), radroots_transport::error::Error> impl<'de> serde_core::de::Deserialize<'de> for radroots_transport::sink::DeliveryRequest pub fn radroots_transport::sink::DeliveryRequest::deserialize<D>(D) -> core::result::Result<Self, <D as serde_core::de::Deserializer>::Error> where D: serde_core::de::Deserializer<'de> pub struct radroots_transport::sink::DeliveryRequestId(_) @@ -304,6 +306,7 @@ pub const fn radroots_transport::SinkStatus::transport_id(&self) -> radroots_tra pub const radroots_transport::sink::DELIVERY_REQUEST_ID_MAX_BYTES: usize pub trait radroots_transport::sink::EventSink: core::marker::Send + core::marker::Sync pub fn radroots_transport::sink::EventSink::deliver(&self, radroots_transport::sink::DeliveryRequest) -> radroots_transport::source::BoxFuture<'_, core::result::Result<radroots_transport::sink::DeliveryReceipt, radroots_transport::sink::SinkFailure>> +pub fn radroots_transport::sink::EventSink::deliver_selected(&self, radroots_transport::sink::DeliveryRequest, radroots_transport::target::TargetSet) -> radroots_transport::source::BoxFuture<'_, core::result::Result<radroots_transport::sink::DeliveryReceipt, radroots_transport::sink::SinkFailure>> pub fn radroots_transport::sink::EventSink::status(&self) -> radroots_transport::source::BoxFuture<'_, core::result::Result<radroots_transport::SinkStatus, radroots_transport::error::Error>> pub mod radroots_transport::source pub enum radroots_transport::source::NextPage @@ -612,6 +615,7 @@ pub radroots_transport::Error::FetchSelectorTooLarge pub radroots_transport::Error::InvalidDeliveryDeadline pub radroots_transport::Error::InvalidDeliveryOutcome pub radroots_transport::Error::InvalidDeliveryRequestId +pub radroots_transport::Error::InvalidDeliveryTargetSelection pub radroots_transport::Error::InvalidDeliveryTimestamp pub radroots_transport::Error::InvalidFetchCursor pub radroots_transport::Error::InvalidFetchDeadline @@ -683,6 +687,7 @@ pub const fn radroots_transport::sink::DeliveryRequest::payload(&self) -> &radro pub const fn radroots_transport::sink::DeliveryRequest::request_id(&self) -> &radroots_transport::sink::DeliveryRequestId pub const fn radroots_transport::sink::DeliveryRequest::satisfaction(&self) -> &radroots_transport::policy::SatisfactionPolicy pub const fn radroots_transport::sink::DeliveryRequest::target_set(&self) -> &radroots_transport::target::TargetSet +pub fn radroots_transport::sink::DeliveryRequest::validate_target_selection(&self, &radroots_transport::target::TargetSet) -> core::result::Result<(), radroots_transport::error::Error> impl<'de> serde_core::de::Deserialize<'de> for radroots_transport::sink::DeliveryRequest pub fn radroots_transport::sink::DeliveryRequest::deserialize<D>(D) -> core::result::Result<Self, <D as serde_core::de::Deserializer>::Error> where D: serde_core::de::Deserializer<'de> pub struct radroots_transport::FetchPage @@ -832,6 +837,7 @@ pub const radroots_transport::TARGET_SET_MAX_ITEMS: usize pub const radroots_transport::TRANSPORT_ID_MAX_BYTES: usize pub trait radroots_transport::EventSink: core::marker::Send + core::marker::Sync pub fn radroots_transport::EventSink::deliver(&self, radroots_transport::sink::DeliveryRequest) -> radroots_transport::source::BoxFuture<'_, core::result::Result<radroots_transport::sink::DeliveryReceipt, radroots_transport::sink::SinkFailure>> +pub fn radroots_transport::EventSink::deliver_selected(&self, radroots_transport::sink::DeliveryRequest, radroots_transport::target::TargetSet) -> radroots_transport::source::BoxFuture<'_, core::result::Result<radroots_transport::sink::DeliveryReceipt, radroots_transport::sink::SinkFailure>> pub fn radroots_transport::EventSink::status(&self) -> radroots_transport::source::BoxFuture<'_, core::result::Result<radroots_transport::SinkStatus, radroots_transport::error::Error>> pub trait radroots_transport::EventSource: core::marker::Send + core::marker::Sync pub fn radroots_transport::EventSource::fetch(&self, radroots_transport::source::FetchRequest) -> radroots_transport::source::BoxFuture<'_, core::result::Result<radroots_transport::source::FetchPage, radroots_transport::error::Error>> diff --git a/contracts/api_baselines/radroots_transport_nostr.txt b/contracts/api_baselines/radroots_transport_nostr.txt @@ -91,10 +91,12 @@ pub fn radroots_transport_nostr::NostrTransport::relay_status(&self) -> radroots impl radroots_transport_nostr::NostrTransport pub fn radroots_transport_nostr::NostrTransport::execute_prepared_delivery(&self, radroots_transport_nostr::PreparedDelivery) -> radroots_transport::source::BoxFuture<'_, core::result::Result<radroots_transport::sink::DeliveryReceipt, radroots_transport::sink::SinkFailure>> pub fn radroots_transport_nostr::NostrTransport::prepare_delivery(&self, radroots_transport::sink::DeliveryRequest) -> core::result::Result<radroots_transport_nostr::PreparedDelivery, alloc::boxed::Box<radroots_transport::sink::SinkFailure>> +pub fn radroots_transport_nostr::NostrTransport::prepare_delivery_selected(&self, radroots_transport::sink::DeliveryRequest, &radroots_transport::target::TargetSet) -> core::result::Result<radroots_transport_nostr::PreparedDelivery, alloc::boxed::Box<radroots_transport::sink::SinkFailure>> impl core::fmt::Debug for radroots_transport_nostr::NostrTransport pub fn radroots_transport_nostr::NostrTransport::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result impl radroots_transport::sink::EventSink for radroots_transport_nostr::NostrTransport pub fn radroots_transport_nostr::NostrTransport::deliver(&self, radroots_transport::sink::DeliveryRequest) -> radroots_transport::source::BoxFuture<'_, core::result::Result<radroots_transport::sink::DeliveryReceipt, radroots_transport::sink::SinkFailure>> +pub fn radroots_transport_nostr::NostrTransport::deliver_selected(&self, radroots_transport::sink::DeliveryRequest, radroots_transport::target::TargetSet) -> radroots_transport::source::BoxFuture<'_, core::result::Result<radroots_transport::sink::DeliveryReceipt, radroots_transport::sink::SinkFailure>> pub fn radroots_transport_nostr::NostrTransport::status(&self) -> radroots_transport::source::BoxFuture<'_, core::result::Result<radroots_transport::status::SinkStatus, radroots_transport::error::Error>> impl radroots_transport::source::EventSource for radroots_transport_nostr::NostrTransport pub fn radroots_transport_nostr::NostrTransport::fetch(&self, radroots_transport::source::FetchRequest) -> radroots_transport::source::BoxFuture<'_, core::result::Result<radroots_transport::source::FetchPage, radroots_transport::error::Error>> diff --git a/contracts/architecture/decisions/delivery_target_selection.v1.json b/contracts/architecture/decisions/delivery_target_selection.v1.json @@ -0,0 +1,14 @@ +{ + "schema": "radroots.delivery-target-selection.v1", + "status": "approved", + "owners": [ + "radroots_transport", + "radroots_transport_nostr", + "radroots_sync", + "radroots_sdk" + ], + "binding": "A constrained attempt passes the exact full frozen DeliveryRequest plus an explicit nonempty bounded TargetSet subset. Selection validates exact Target equality before new claims or transport I/O. It never changes signed event bytes, request identity, satisfaction policy, deadline or persisted request. No subset request or receipt rebinding is permitted.", + "adapter": "EventSink deliver_selected defaults to ordinary delivery only for all original targets; a proper subset is unsupported without I/O. Nostr preparation retains the full request and withholds all unselected endpoints before configuration admission with unattempted retryable target_not_selected rows. Existing configuration denial and backoff remain additional constraints. SDK NostrSlot forwards the selection to the retained adapter snapshot without fallback. SDK Sync Operations forwards selected delivery once to the canonical engine, retaining the caller subset and original durable request.", + "evidence": "Successful receipts remain complete and ordered against the original target set. Partial failures retain their original full request binding. An attempted row outside the selection is an invalid adapter contract, never fabricated success. Shared Sync claims, raw fact persistence, late-result reconciliation, stop and retry bounds remain unchanged. Selection is caller scheduling input, not a second durable request or journal.", + "scope": "Host product policy determines eligibility and holds empty selections without invoking delivery. Shared libraries do not know replacement or deletion policy. Ordinary unconstrained delivery, persistence schemas, dependencies, coverage thresholds and release authorization remain unchanged." +} diff --git a/crates/sdk/src/sync.rs b/crates/sdk/src/sync.rs @@ -206,6 +206,17 @@ impl<'a> Operations<'a> { self.engine.deliver_push(operation_id).await } + /// Attempts an exact subset while retaining the full durable request binding. + pub async fn deliver_push_selected( + &self, + operation_id: radroots_sync::policy::SyncId, + selected: radroots_transport::TargetSet, + ) -> Result<DeliveryExecutionReceipt, Error> { + self.engine + .deliver_push_selected(operation_id, selected) + .await + } + /// Returns the native passive sync status without starting recovery work. pub async fn status(&self, projections: &[ProjectionId]) -> Result<SyncStatus, Error> { self.engine.status(projections).await @@ -233,6 +244,8 @@ impl std::fmt::Debug for Operations<'_> { #[cfg(all(test, feature = "sync", feature = "memory"))] mod tests { + #[cfg(feature = "local-signing")] + mod selected; use std::sync::{ Arc, atomic::{AtomicU8, Ordering}, diff --git a/crates/sdk/src/sync/tests/selected.rs b/crates/sdk/src/sync/tests/selected.rs @@ -0,0 +1,150 @@ +use super::*; +use radroots_transport::{ + DeliveryReceipt, DeliveryRequest, EventSink, SinkFailure, + outcome::DeliveryOutcome, + sink::{DeliveryTargetReceipt, SinkStatus}, +}; + +#[derive(Default)] +struct SelectedSink(std::sync::Mutex<Vec<(DeliveryRequest, TargetSet)>>); + +impl EventSink for SelectedSink { + fn status(&self) -> radroots_transport::BoxFuture<'_, Result<SinkStatus, TransportError>> { + Box::pin(async { panic!("selected delivery does not probe status") }) + } + + fn deliver( + &self, + _: DeliveryRequest, + ) -> radroots_transport::BoxFuture<'_, Result<DeliveryReceipt, SinkFailure>> { + Box::pin(async { panic!("selected delivery must not widen to ordinary delivery") }) + } + + fn deliver_selected( + &self, + request: DeliveryRequest, + selected: TargetSet, + ) -> radroots_transport::BoxFuture<'_, Result<DeliveryReceipt, SinkFailure>> { + Box::pin(async move { + request.validate_target_selection(&selected).unwrap(); + self.0 + .lock() + .unwrap() + .push((request.clone(), selected.clone())); + Ok(DeliveryReceipt::for_request( + &request, + request + .target_set() + .targets() + .iter() + .map(|target| { + if selected.targets().contains(target) { + DeliveryTargetReceipt::attempted( + target.clone(), + DeliveryOutcome::accepted(), + ) + } else { + DeliveryTargetReceipt::skipped( + target.clone(), + DeliveryOutcome::unavailable(), + ) + .unwrap() + } + }) + .collect(), + ) + .unwrap()) + }) + } +} + +#[tokio::test] +async fn sdk_selected_delivery_preserves_subset_and_full_durable_request() { + let storage = Arc::new(MemoryStorage::new( + SourceGeneration::new([211; 32]).unwrap(), + )); + let signer = radroots_nostr::signing::LocalSigner::new( + radroots_nostr::key::SecretKey::parse( + "0000000000000000000000000000000000000000000000000000000000000001", + ) + .unwrap(), + ) + .unwrap(); + let author = "79be667ef9dcbbac55a06295ce870b07029bfcdb2dce28d959f2815b16f81798"; + let (clock, ids, deadlines) = HostPolicy::default().composition(); + let now = clock.now_unix_ms().unwrap(); + let sink = Arc::new(SelectedSink::default()); + let engine = Engine::builder(storage.clone(), clock, ids, deadlines) + .signer(Arc::new(signer)) + .sink(sink.clone()) + .build() + .unwrap(); + let client = ClientBuilder::new() + .storage(storage) + .sync_engine(engine) + .build() + .unwrap(); + let operations = client.sync().unwrap().unwrap(); + let targets = TargetSet::new(vec![ + target(), + Target::nostr_relay("wss://held.example").unwrap(), + ]) + .unwrap(); + let selected = TargetSet::new(vec![targets.targets()[0].clone()]).unwrap(); + let request = PushRequest::new( + SyncId::new([211; 16]).unwrap(), + IdempotencyKey::parse("sdk-selected-delivery").unwrap(), + Actor::new( + PublicKey::from_hex(author).unwrap(), + ActorSource::ExplicitPublicKey, + [AuthorRole::Any], + ) + .unwrap(), + AuthoredEventPlan::from_generic( + GenericEventDraft::new( + "radroots.social.geochat.v1", + 20_000, + now / 1_000, + Vec::new(), + "selected", + author, + ) + .unwrap(), + ) + .unwrap(), + targets.clone(), + SatisfactionPolicy::new(SatisfactionClass::Accepted, TargetPolicy::all()), + now + 60_000, + CancellationPolicy::LocalCooperative, + ) + .unwrap(); + let id = request.operation_id(); + operations.submit_push(request).await.unwrap(); + let original = operations.push_status(id).await.unwrap().unwrap(); + let original_request = original.delivery_plan().request().unwrap(); + assert_eq!( + operations + .deliver_push_selected( + id, + TargetSet::new(vec![Target::nostr_relay("wss://foreign.example").unwrap()]) + .unwrap(), + ) + .await + .unwrap_err(), + Error::InvalidDeliveryRequest + ); + assert!(sink.0.lock().unwrap().is_empty()); + operations + .deliver_push_selected(id, selected.clone()) + .await + .unwrap(); + let after = operations.push_status(id).await.unwrap().unwrap(); + assert_eq!(after.delivery_plan().request(), Some(original_request)); + assert_eq!( + sink.0.lock().unwrap().as_slice(), + &[(original_request.clone(), selected)] + ); + assert!(!after.delivery_plan().state().is_terminal()); + assert_eq!(after.delivery_plan().delivery_facts().len(), 1); + client.close().await.unwrap(); +} diff --git a/crates/sdk/src/transport.rs b/crates/sdk/src/transport.rs @@ -2052,6 +2052,38 @@ impl radroots_transport::EventSink for NostrSlot { radroots_transport::EventSink::deliver(state.transport.as_ref(), request).await }) } + + fn deliver_selected( + &self, + request: radroots_transport::DeliveryRequest, + selected: radroots_transport::TargetSet, + ) -> radroots_transport::BoxFuture< + '_, + Result<radroots_transport::DeliveryReceipt, radroots_transport::SinkFailure>, + > { + Box::pin(async move { + if request.validate_target_selection(&selected).is_err() { + return Err(radroots_transport::SinkFailure::invalid_contract(&request)); + } + let Some(state) = self.snapshot() else { + return Err(radroots_transport::SinkFailure::for_request( + &request, + "nostr_transport_not_configured", + radroots_transport::outcome::Retryability::Terminal, + None, + None, + Vec::new(), + ) + .expect("static unconfigured sink failure is valid")); + }; + radroots_transport::EventSink::deliver_selected( + state.transport.as_ref(), + request, + selected, + ) + .await + }) + } } #[cfg(feature = "nostr")] @@ -3420,7 +3452,66 @@ mod tests { 1, ) .expect("delivery"); - let failure = slot.deliver(deliver).await.expect_err("unconfigured sink"); + let failure = slot + .deliver(deliver.clone()) + .await + .expect_err("unconfigured sink"); + assert_eq!(failure.code(), "nostr_transport_not_configured"); + let failure = slot + .deliver_selected(deliver.clone(), deliver.target_set().clone()) + .await + .unwrap_err(); assert_eq!(failure.code(), "nostr_transport_not_configured"); + failure.validate_for_request(&deliver).unwrap(); + let foreign = TargetSet::new(vec![target(2)]).unwrap(); + let failure = slot + .deliver_selected(deliver.clone(), foreign) + .await + .unwrap_err(); + assert_eq!(failure.code(), "invalid_transport_contract"); + let subset_request = DeliveryRequest::new( + "selected-delivery", + DeliveryPayload::new(signed_event()), + TargetSet::new(vec![target(1), target(2)]).unwrap(), + SatisfactionPolicy::new(SatisfactionClass::Accepted, TargetPolicy::all()), + 1, + ) + .unwrap(); + slot.configure( + RelayProfile::explicit( + RelayProfileKind::Public, + [RelayEndpoint::new( + "wss://one.example", + RelayUrlPolicy::Public, + RelayAccess::ReadWrite, + ) + .unwrap()], + ) + .unwrap(), + ) + .unwrap(); + let receipt = slot + .deliver_selected(deliver.clone(), deliver.target_set().clone()) + .await + .unwrap(); + receipt.validate_for_request(&deliver).unwrap(); + assert!( + receipt + .target_receipts() + .iter() + .all(|row| !row.was_attempted()) + ); + let subset = + TargetSet::new(vec![subset_request.target_set().targets()[0].clone()]).unwrap(); + let receipt = slot + .deliver_selected(subset_request.clone(), subset) + .await + .unwrap(); + receipt.validate_for_request(&subset_request).unwrap(); + assert!(!receipt.target_receipts()[1].was_attempted()); + assert_eq!( + receipt.target_receipts()[1].outcome().code(), + Some("target_not_selected") + ); } } diff --git a/crates/sdk/tests/package_boundary.rs b/crates/sdk/tests/package_boundary.rs @@ -173,6 +173,7 @@ fn package_contains_only_reachable_sources_and_registered_targets() { "signing.rs".to_owned(), "storage.rs".to_owned(), "sync.rs".to_owned(), + "sync/tests/selected.rs".to_owned(), "trade.rs".to_owned(), "transport.rs".to_owned(), ]) diff --git a/crates/sync/README.md b/crates/sync/README.md @@ -61,3 +61,10 @@ After a post-delivery clock failure, Sync retains the raw non-expiring result with the known pre-effect time as a causal lower bound and returns `ClockUnavailable` without retry scheduling. This is not a measured response time and does not change strict signing or expiring authorization requirements. + +`Engine::deliver_push_selected` accepts an explicit nonempty subset of the +frozen delivery targets. It validates the subset before a new claim and passes +the original full request to the selected sink boundary. Attempted evidence +outside that selection is an invalid adapter contract. Shared claim, stop, +late-result persistence, reconciliation and retry semantics remain unchanged; +the caller owns eligibility and holds an empty selection without delivery. diff --git a/crates/sync/src/push/delivery.rs b/crates/sync/src/push/delivery.rs @@ -13,6 +13,25 @@ impl Engine { &self, operation_id: SyncId, ) -> Result<DeliveryExecutionReceipt, Error> { + self.deliver_push_inner(operation_id, None).await + } + + /// Attempts only an exact nonempty subset of the frozen delivery targets. + /// The full persisted request, claim and raw result bindings remain intact. + /// Callers own selection policy; no ineligible target may be attempted. + pub async fn deliver_push_selected( + &self, + operation_id: SyncId, + selected: radroots_transport::TargetSet, + ) -> Result<DeliveryExecutionReceipt, Error> { + self.deliver_push_inner(operation_id, Some(selected)).await + } + + async fn deliver_push_inner( + &self, + operation_id: SyncId, + selected: Option<radroots_transport::TargetSet>, + ) -> Result<DeliveryExecutionReceipt, Error> { let status = self.push_status(operation_id).await?.ok_or_else(|| { if self.sink.is_none() { Error::MissingSink @@ -56,6 +75,11 @@ impl Engine { } let sink = self.sink.as_deref().ok_or(Error::MissingSink)?; let request = plan.request().cloned().ok_or(Error::InvalidSignerOutput)?; + if let Some(selected) = &selected { + request + .validate_target_selection(selected) + .map_err(|_| Error::InvalidDeliveryRequest)?; + } let claimed = self.claim_delivery_plan(plan, now).await?; let claim = claimed .claim_evidence() @@ -99,11 +123,27 @@ impl Engine { .map_err(|_| Error::InvalidDeliveryRequest)?, ) } else { - match sink.deliver(request.clone()).await { - Ok(receipt) if receipt.validate_for_request(&request).is_ok() => { + let result = match selected.clone() { + Some(targets) => sink.deliver_selected(request.clone(), targets).await, + None => sink.deliver(request.clone()).await, + }; + let allowed = |rows: &[radroots_transport::sink::DeliveryTargetReceipt]| { + selected.as_ref().is_none_or(|targets| { + rows.iter() + .all(|row| !row.was_attempted() || targets.targets().contains(row.target())) + }) + }; + match result { + Ok(receipt) + if receipt.validate_for_request(&request).is_ok() + && allowed(receipt.target_receipts()) => + { DeliveryAttemptOutcome::Receipt(receipt) } - Err(failure) if failure.validate_for_request(&request).is_ok() => { + Err(failure) + if failure.validate_for_request(&request).is_ok() + && allowed(failure.partial_evidence()) => + { DeliveryAttemptOutcome::SinkFailure(failure) } Ok(_) | Err(_) => { diff --git a/crates/sync/tests/push_enqueue.rs b/crates/sync/tests/push_enqueue.rs @@ -60,6 +60,9 @@ mod signing_evidence; #[path = "push_enqueue/delivery_evidence.rs"] mod delivery_evidence; +#[path = "push_enqueue/delivery_selection.rs"] +mod delivery_selection; + struct MockSink; struct FaultStorage { diff --git a/crates/sync/tests/push_enqueue/delivery_evidence.rs b/crates/sync/tests/push_enqueue/delivery_evidence.rs @@ -58,7 +58,7 @@ fn setup( execute_to_admitted(&engine, &push); (engine, storage, clock, sink, push) } -fn source_only(storage: Arc<dyn SyncStorage>, clock: Arc<dyn Clock>) -> Engine { +pub(super) fn source_only(storage: Arc<dyn SyncStorage>, clock: Arc<dyn Clock>) -> Engine { Engine::builder( storage, clock, diff --git a/crates/sync/tests/push_enqueue/delivery_selection.rs b/crates/sync/tests/push_enqueue/delivery_selection.rs @@ -0,0 +1,308 @@ +use super::*; +use futures::channel::oneshot; +use radroots_storage::authored_delivery::DeliveryAttemptOutcome; +use radroots_transport::policy::SatisfactionState; + +type Pending = ( + DeliveryRequest, + TargetSet, + oneshot::Sender<Result<DeliveryReceipt, SinkFailure>>, +); + +#[derive(Default)] +struct SelectedSink { + pending: Mutex<VecDeque<Pending>>, + calls: AtomicUsize, +} + +impl EventSink for SelectedSink { + fn status(&self) -> radroots_transport::BoxFuture<'_, Result<SinkStatus, TransportError>> { + Box::pin(async { panic!("no status probe") }) + } + fn deliver( + &self, + _: DeliveryRequest, + ) -> radroots_transport::BoxFuture<'_, Result<DeliveryReceipt, SinkFailure>> { + Box::pin(async { panic!("selected attempts must not fall back") }) + } + fn deliver_selected( + &self, + request: DeliveryRequest, + selected: TargetSet, + ) -> radroots_transport::BoxFuture<'_, Result<DeliveryReceipt, SinkFailure>> { + Box::pin(async move { + request.validate_target_selection(&selected).unwrap(); + self.calls.fetch_add(1, Ordering::SeqCst); + let (send, receive) = oneshot::channel(); + self.pending + .lock() + .unwrap() + .push_back((request, selected, send)); + receive.await.unwrap() + }) + } +} + +fn selected_receipt(request: &DeliveryRequest, selected: &TargetSet) -> DeliveryReceipt { + DeliveryReceipt::for_request( + request, + request + .target_set() + .targets() + .iter() + .map(|target| { + if selected.targets().contains(target) { + DeliveryTargetReceipt::attempted(target.clone(), DeliveryOutcome::accepted()) + } else { + DeliveryTargetReceipt::skipped(target.clone(), DeliveryOutcome::unavailable()) + .unwrap() + } + }) + .collect(), + ) + .unwrap() +} + +fn poll_pending(future: &mut (impl std::future::Future + Unpin)) { + assert!( + future + .poll_unpin(&mut std::task::Context::from_waker(noop_waker_ref())) + .is_pending() + ); +} + +fn engine(storage: Arc<dyn SyncStorage>, clock: Arc<TestClock>, sink: Arc<SelectedSink>) -> Engine { + Engine::builder( + storage, + clock, + Arc::new(TestIds(AtomicU64::new(10))), + DeadlinePolicy::new(10_000, 10_000, 10_000).unwrap(), + ) + .sink(sink) + .signer(Arc::new(MockSigner::new(SignBehavior::Success { + completed_at_unix_ms: 1_800_000_200_500, + }))) + .build() + .unwrap() +} + +fn push() -> PushRequest { + request_with_policy( + 211, + &["wss://one.example", "wss://two.example"], + SatisfactionClass::Accepted, + TargetPolicy::all(), + ) +} + +#[test] +fn invalid_selection_cannot_claim_and_out_of_selection_results_fail_closed() { + for failure in [false, true] { + let storage = Arc::new(MemoryStorage::new( + SourceGeneration::new([211; 32]).unwrap(), + )); + let clock = Arc::new(TestClock(AtomicU64::new(1_800_000_200_000))); + let sink = Arc::new(SelectedSink::default()); + let engine = engine(storage, clock, sink.clone()); + let push = push(); + execute_to_admitted(&engine, &push); + let before = block_on(engine.push_status(push.operation_id())) + .unwrap() + .unwrap(); + let original = before.delivery_plan().request().unwrap(); + let foreign = + TargetSet::new(vec![Target::nostr_relay("wss://foreign.example").unwrap()]).unwrap(); + assert!(matches!( + block_on(engine.deliver_push_selected(push.operation_id(), foreign)), + Err(Error::InvalidDeliveryRequest) + )); + let unchanged = block_on(engine.push_status(push.operation_id())) + .unwrap() + .unwrap(); + assert_eq!(unchanged.delivery_plan(), before.delivery_plan()); + assert_eq!(sink.calls.load(Ordering::SeqCst), 0); + let selected = TargetSet::new(vec![original.target_set().targets()[0].clone()]).unwrap(); + let mut future = Box::pin(engine.deliver_push_selected(push.operation_id(), selected)); + poll_pending(&mut future); + let (request, _, sender) = sink.pending.lock().unwrap().pop_front().unwrap(); + assert_eq!(&request, original); + let forbidden = DeliveryTargetReceipt::attempted( + request.target_set().targets()[1].clone(), + DeliveryOutcome::accepted(), + ); + let result = if failure { + Err(SinkFailure::for_request( + &request, + "upstream_failure", + Retryability::Retryable, + None, + None, + vec![forbidden], + ) + .unwrap()) + } else { + Ok(receipt(&request, vec![DeliveryOutcome::accepted(); 2]).unwrap()) + }; + sender.send(result).unwrap(); + let result = block_on(future).unwrap(); + assert_eq!(result.plan().request(), Some(original)); + let DeliveryAttemptOutcome::SinkFailure(failure) = + result.plan().delivery_facts()[0].outcome() + else { + panic!("invalid adapter evidence is never acceptance") + }; + assert_eq!(failure.code(), "invalid_transport_contract"); + assert!(failure.partial_evidence().is_empty()); + assert_ne!( + result.plan().delivery_satisfaction().unwrap(), + SatisfactionState::Satisfied + ); + } +} + +#[test] +fn selected_late_acceptance_survives_stop_without_scheduling_new_targets() { + let storage = Arc::new(MemoryStorage::new( + SourceGeneration::new([212; 32]).unwrap(), + )); + let clock = Arc::new(TestClock(AtomicU64::new(1_800_000_200_000))); + let sink = Arc::new(SelectedSink::default()); + let engine = engine(storage.clone(), clock.clone(), sink.clone()); + let push = push(); + execute_to_admitted(&engine, &push); + let before = block_on(engine.push_status(push.operation_id())) + .unwrap() + .unwrap(); + let request = before.delivery_plan().request().unwrap(); + let selected = TargetSet::new(vec![request.target_set().targets()[0].clone()]).unwrap(); + let mut future = Box::pin(engine.deliver_push_selected(push.operation_id(), selected.clone())); + poll_pending(&mut future); + block_on(engine.cancel_push(push.operation_id())).unwrap(); + let (captured, captured_selection, sender) = sink.pending.lock().unwrap().pop_front().unwrap(); + assert_eq!(&captured, request); + assert_eq!(captured_selection, selected); + let expected = selected_receipt(&captured, &selected); + sender.send(Ok(expected.clone())).unwrap(); + let late = block_on(future).unwrap(); + assert_eq!(late.plan().state(), AuthoredDeliveryState::Cancelled); + assert_eq!(late.plan().attempt_count(), 0); + assert_eq!( + late.plan().delivery_facts()[0].outcome(), + &DeliveryAttemptOutcome::Receipt(expected) + ); + let recovery = delivery_evidence::source_only(storage, clock); + let replay = block_on(recovery.deliver_push_selected(push.operation_id(), selected)).unwrap(); + assert!(replay.is_replay()); + assert_eq!(replay.plan(), late.plan()); + assert_eq!(sink.calls.load(Ordering::SeqCst), 1); +} + +#[tokio::test] +async fn sqlite_selected_facts_survive_reopen_and_next_target_finishes_same_request() { + let directory = tempfile::tempdir().unwrap(); + let paths = Paths::from_directory(directory.path()).unwrap(); + let storage = Arc::new( + SqliteStorage::open( + OpenOptions::new(paths.clone(), OpenMode::Create) + .with_source_generation(SourceGeneration::new([213; 32]).unwrap(), 1) + .unwrap(), + ) + .await + .unwrap(), + ); + let clock = Arc::new(TestClock(AtomicU64::new(1_800_000_200_000))); + let sink = Arc::new(SelectedSink::default()); + let first_engine = engine(storage.clone(), clock.clone(), sink.clone()); + let push = push(); + first_engine.sign_prepared(push.clone()).await.unwrap(); + first_engine + .admit_signed(push.operation_id()) + .await + .unwrap(); + let before = first_engine + .push_status(push.operation_id()) + .await + .unwrap() + .unwrap(); + let original = before.delivery_plan().request().unwrap().clone(); + let selected_a = TargetSet::new(vec![original.target_set().targets()[0].clone()]).unwrap(); + let first = { + let mut future = + Box::pin(first_engine.deliver_push_selected(push.operation_id(), selected_a.clone())); + // SQLite storage needs an executor turn before the sink is reached. + let responder = async { + let pending = loop { + if let Some(pending) = sink.pending.lock().unwrap().pop_front() { + break pending; + } + tokio::task::yield_now().await; + }; + assert_eq!(pending.0, original); + assert_eq!(pending.1, selected_a); + pending + .2 + .send(Ok(selected_receipt(&pending.0, &pending.1))) + .unwrap(); + }; + tokio::time::timeout(std::time::Duration::from_secs(10), async { + let (result, ()) = tokio::join!(&mut future, responder); + result.unwrap() + }) + .await + .unwrap() + }; + assert_ne!(first.plan().state(), AuthoredDeliveryState::Satisfied); + assert_eq!(first.plan().attempt_count(), 1); + clock.0.store( + first.plan().retry().unwrap().not_before_unix_ms(), + Ordering::SeqCst, + ); + drop(first_engine); + storage.close().await.unwrap(); + let storage = Arc::new( + SqliteStorage::open(OpenOptions::new(paths, OpenMode::ReadWriteExisting)) + .await + .unwrap(), + ); + let second_engine = engine(storage.clone(), clock.clone(), sink.clone()); + let recovered = second_engine + .push_status(push.operation_id()) + .await + .unwrap() + .unwrap(); + assert_eq!( + recovered.delivery_plan().delivery_facts(), + first.plan().delivery_facts() + ); + let selected_b = TargetSet::new(vec![original.target_set().targets()[1].clone()]).unwrap(); + let responder = async { + let pending = loop { + if let Some(pending) = sink.pending.lock().unwrap().pop_front() { + break pending; + } + tokio::task::yield_now().await; + }; + assert_eq!(pending.0, original); + assert_eq!(pending.1, selected_b); + pending + .2 + .send(Ok(selected_receipt(&pending.0, &pending.1))) + .unwrap(); + }; + let finished = tokio::time::timeout(std::time::Duration::from_secs(10), async { + let (result, ()) = tokio::join!( + second_engine.deliver_push_selected(push.operation_id(), selected_b.clone()), + responder + ); + result.unwrap() + }) + .await + .unwrap(); + assert_eq!(finished.plan().request(), Some(&original)); + assert_eq!(finished.plan().state(), AuthoredDeliveryState::Satisfied); + assert_eq!(finished.plan().attempt_count(), 2); + assert_eq!(finished.plan().delivery_facts().len(), 2); + assert_eq!(sink.calls.load(Ordering::SeqCst), 2); + drop(second_engine); + storage.close().await.unwrap(); +} diff --git a/crates/transport/README.md b/crates/transport/README.md @@ -78,6 +78,12 @@ does not select an async runtime or require an async-trait macro. Native sources may be retained by adapters, but public outcome codes and messages must remain bounded and secret-safe. +`EventSink::deliver_selected` accepts an exact nonempty subset while preserving +the full original request and receipt target set. An adapter must not attempt +unselected targets. The default accepts only all original targets; a proper +subset returns `target_selection_unsupported` without delivery. Callers hold +empty selections without invoking transport. + ## Targets and extensible identity `TransportId` is a validated open identity, not a closed enum. The built-in diff --git a/crates/transport/src/error.rs b/crates/transport/src/error.rs @@ -59,6 +59,7 @@ pub enum Error { InvalidDeliveryRequestId, InvalidDeliveryTimestamp, InvalidDeliveryDeadline, + InvalidDeliveryTargetSelection, InvalidDeliveryOutcome, UnexpectedDeliveryTargetReceipt, DuplicateDeliveryTargetReceipt, @@ -197,6 +198,9 @@ impl fmt::Display for Error { f.write_str("transport delivery timestamp is invalid") } Self::InvalidDeliveryDeadline => f.write_str("transport delivery deadline is invalid"), + Self::InvalidDeliveryTargetSelection => { + f.write_str("transport delivery target selection is not an exact subset") + } Self::InvalidDeliveryOutcome => f.write_str("transport delivery outcome is invalid"), Self::UnexpectedDeliveryTargetReceipt => { f.write_str("transport delivery receipt contains an unexpected target") diff --git a/crates/transport/src/sink.rs b/crates/transport/src/sink.rs @@ -8,6 +8,7 @@ use crate::{ target::{Target, TargetSet}, }; use alloc::{ + boxed::Box, collections::{BTreeMap, BTreeSet}, string::{String, ToString}, vec::Vec, @@ -135,6 +136,20 @@ impl DeliveryRequest { pub const fn deadline_unix_ms(&self) -> u64 { self.deadline_unix_ms } + + /// Validates a nonempty bounded subset of the exact original targets. + /// Selection never changes this request's payload, policy or identity. + pub fn validate_target_selection(&self, selected: &TargetSet) -> Result<(), Error> { + if selected + .targets() + .iter() + .all(|target| self.target_set.targets().contains(target)) + { + Ok(()) + } else { + Err(Error::InvalidDeliveryTargetSelection) + } + } } /// Normalized result for one requested target. @@ -418,6 +433,36 @@ pub trait EventSink: Send + Sync { &self, request: DeliveryRequest, ) -> BoxFuture<'_, Result<DeliveryReceipt, SinkFailure>>; + + /// Attempts only selected targets while retaining the full request binding. + /// + /// Implementations must validate the exact subset before I/O, keep every + /// original target in successful receipts, and report unselected targets as + /// unattempted. The default supports only a selection of every target; a + /// proper subset fails closed without calling ordinary delivery. + fn deliver_selected( + &self, + request: DeliveryRequest, + selected: TargetSet, + ) -> BoxFuture<'_, Result<DeliveryReceipt, SinkFailure>> { + Box::pin(async move { + if request.validate_target_selection(&selected).is_err() { + return Err(SinkFailure::invalid_contract(&request)); + } + if selected.len() == request.target_set().len() { + return self.deliver(request).await; + } + Err(SinkFailure::for_request( + &request, + "target_selection_unsupported", + Retryability::Terminal, + None, + None, + Vec::new(), + ) + .expect("static unsupported selection failure is valid")) + }) + } } #[cfg(feature = "serde")] diff --git a/crates/transport/tests/delivery_contract.rs b/crates/transport/tests/delivery_contract.rs @@ -56,6 +56,78 @@ fn mixed_receipt(request: &DeliveryRequest) -> DeliveryReceipt { } #[test] +fn target_selection_is_exact_and_default_adapter_never_widens_it() { + use radroots_transport::{BoxFuture, EventSink, SinkStatus}; + use std::sync::atomic::{AtomicUsize, Ordering}; + struct Sink(AtomicUsize); + impl EventSink for Sink { + fn status(&self) -> BoxFuture<'_, Result<SinkStatus, Error>> { + Box::pin(async { panic!("no status probe") }) + } + fn deliver( + &self, + request: DeliveryRequest, + ) -> BoxFuture<'_, Result<DeliveryReceipt, SinkFailure>> { + Box::pin(async move { + self.0.fetch_add(1, Ordering::SeqCst); + Ok(mixed_receipt(&request)) + }) + } + } + let sink = Sink(AtomicUsize::new(0)); + let request = request(SatisfactionPolicy::new( + SatisfactionClass::Accepted, + TargetPolicy::all(), + )); + let original = request.clone(); + let labelled = Target::new_with_metadata( + radroots_transport::TransportId::NOSTR, + "wss://one.example", + None, + Some(radroots_transport::target::TargetLabel::parse("changed label").unwrap()), + ) + .unwrap(); + assert_eq!( + labelled.fingerprint(), + request.target_set().targets()[0].fingerprint() + ); + assert_eq!( + request.validate_target_selection(&TargetSet::new(vec![labelled]).unwrap()), + Err(Error::InvalidDeliveryTargetSelection) + ); + let subset = TargetSet::new(vec![request.target_set().targets()[0].clone()]).unwrap(); + request.validate_target_selection(&subset).unwrap(); + let unsupported = + futures::executor::block_on(sink.deliver_selected(request.clone(), subset)).unwrap_err(); + unsupported.validate_for_request(&request).unwrap(); + assert_eq!(unsupported.code(), "target_selection_unsupported"); + let foreign = + TargetSet::new(vec![Target::nostr_relay("wss://foreign.example").unwrap()]).unwrap(); + assert_eq!( + request.validate_target_selection(&foreign), + Err(Error::InvalidDeliveryTargetSelection) + ); + assert_eq!( + Error::InvalidDeliveryTargetSelection.to_string(), + "transport delivery target selection is not an exact subset" + ); + let invalid = + futures::executor::block_on(sink.deliver_selected(request.clone(), foreign)).unwrap_err(); + invalid.validate_for_request(&request).unwrap(); + assert_eq!(invalid.code(), "invalid_transport_contract"); + assert_eq!(sink.0.load(Ordering::SeqCst), 0); + let mut reversed = request.target_set().targets().to_vec(); + reversed.reverse(); + let receipt = futures::executor::block_on( + sink.deliver_selected(request.clone(), TargetSet::new(reversed).unwrap()), + ) + .unwrap(); + receipt.validate_for_request(&request).unwrap(); + assert_eq!(sink.0.load(Ordering::SeqCst), 1); + assert_eq!(request, original); +} + +#[test] fn any_all_quorum_and_required_targets_are_exact() { let any = request(SatisfactionPolicy::new( SatisfactionClass::Accepted, diff --git a/crates/transport/tests/source_contract.rs b/crates/transport/tests/source_contract.rs @@ -47,21 +47,24 @@ fn fetch_selector_is_bounded_canonical_and_request_bound() { assert_eq!(exact_tags[0].0, 'd'); assert_eq!(exact_tags[0].1, &[String::from("trade-1")]); assert!(selector.matches(&event)); - let encoded = serde_json::to_string(&selector).expect("selector JSON"); - assert_eq!( - serde_json::from_str::<FetchSelector>(encoded.as_str()).expect("selector round trip"), - selector - ); - assert!( - serde_json::from_value::<FetchSelector>(serde_json::json!({ - "kinds": [], - "authors": [], - "exact_tags": {"D": ["trade-1"]}, - "since_unix_seconds": null, - "until_unix_seconds": null - })) - .is_err() - ); + #[cfg(feature = "serde")] + { + let encoded = serde_json::to_string(&selector).expect("selector JSON"); + assert_eq!( + serde_json::from_str::<FetchSelector>(encoded.as_str()).expect("selector round trip"), + selector + ); + assert!( + serde_json::from_value::<FetchSelector>(serde_json::json!({ + "kinds": [], + "authors": [], + "exact_tags": {"D": ["trade-1"]}, + "since_unix_seconds": null, + "until_unix_seconds": null + })) + .is_err() + ); + } assert_eq!( FetchSelector::all() .with_kinds(vec![1, 1]) @@ -147,6 +150,7 @@ fn fetch_selector_is_bounded_canonical_and_request_bound() { .count(), radroots_transport::source::FETCH_SELECTOR_MAX_TAG_KEYS ); + #[cfg(feature = "serde")] assert!( serde_json::from_str::<FetchSelector>( r#"{"kinds":[],"authors":[],"exact_tags":{"d":["one"],"d":["two"]},"since_unix_seconds":null,"until_unix_seconds":null}"#, @@ -183,7 +187,10 @@ fn tagged_event() -> SignedEvent { let mut wire = Nip01EventWire::parse_json(raw).expect("wire event"); wire.tags = vec![vec![String::from("d"), String::from("trade-1")]]; wire.id = wire.computed_event_id().expect("event id").into_string(); - let raw = serde_json::to_string(&wire).expect("event JSON"); + let mut value: serde_json::Value = serde_json::from_str(raw).expect("fixture JSON"); + value["id"] = serde_json::json!(&wire.id); + value["tags"] = serde_json::json!(&wire.tags); + let raw = value.to_string(); SignedEvent::from_wire_verified_id(wire, raw.as_str()).expect("signed event") } diff --git a/crates/transport_nostr/README.md b/crates/transport_nostr/README.md @@ -89,6 +89,13 @@ let _forged = PreparedDelivery { }; ``` +`prepare_delivery_selected` and `EventSink::deliver_selected` further constrain +an attempt to an exact subset of the original targets. The retained request, +signed bytes and full receipt binding do not change. Unselected targets receive +unattempted, retryable `target_not_selected` rows before configuration or I/O +admission. Selected targets still require writable configuration and backoff +admission. This adapter does not choose application eligibility policy. + ## Public surface - [`RelayProfile`] defines public, loopback-simulator, and physical-device diff --git a/crates/transport_nostr/src/sink.rs b/crates/transport_nostr/src/sink.rs @@ -5,7 +5,7 @@ use crate::{NostrTransport, RelayUrl, status}; use core::{fmt, time::Duration}; use futures::{StreamExt, stream}; use radroots_transport::{ - BoxFuture, DeliveryReceipt, DeliveryRequest, EventSink, SinkFailure, Target, + BoxFuture, DeliveryReceipt, DeliveryRequest, EventSink, SinkFailure, Target, TargetSet, outcome::DeliveryOutcome, sink::{DeliveryTargetReceipt, SinkStatus}, }; @@ -142,9 +142,45 @@ impl NostrTransport { &self, request: DeliveryRequest, ) -> Result<PreparedDelivery, Box<SinkFailure>> { + self.prepare_delivery_inner(request, None) + } + + /// Prepares only an exact subset, retaining the complete original request. + /// Unselected targets remain explicit unattempted, retryable receipt rows. + pub fn prepare_delivery_selected( + &self, + request: DeliveryRequest, + selected: &TargetSet, + ) -> Result<PreparedDelivery, Box<SinkFailure>> { + request + .validate_target_selection(selected) + .map_err(|_| Box::new(SinkFailure::invalid_contract(&request)))?; + self.prepare_delivery_inner(request, Some(selected)) + } + + fn prepare_delivery_inner( + &self, + request: DeliveryRequest, + selected: Option<&TargetSet>, + ) -> Result<PreparedDelivery, Box<SinkFailure>> { let mut authorized = Vec::new(); let mut skipped = Vec::new(); for target in request.target_set().targets() { + if selected.is_some_and(|targets| !targets.targets().contains(target)) { + skipped.push( + DeliveryTargetReceipt::skipped( + target.clone(), + DeliveryOutcome::unavailable() + .with_detail( + "target_not_selected", + "target is held for a later attempt", + ) + .map_err(|_| Box::new(SinkFailure::invalid_contract(&request)))?, + ) + .map_err(|_| Box::new(SinkFailure::invalid_contract(&request)))?, + ); + continue; + } match self.config().endpoint_for_target(target) { Some(endpoint) if endpoint.access().can_write() => { authorized.push((endpoint.url().clone(), target.clone())); @@ -293,6 +329,19 @@ impl EventSink for NostrTransport { self.execute_prepared_delivery(prepared).await }) } + + fn deliver_selected( + &self, + request: DeliveryRequest, + selected: TargetSet, + ) -> BoxFuture<'_, Result<DeliveryReceipt, SinkFailure>> { + Box::pin(async move { + let prepared = self + .prepare_delivery_selected(request, &selected) + .map_err(|failure| *failure)?; + self.execute_prepared_delivery(prepared).await + }) + } } #[cfg_attr(coverage_nightly, coverage(off))] @@ -422,6 +471,59 @@ mod tests { } #[test] + fn selected_preparation_retains_binding_and_only_authorizes_selected_relays() { + let config = Config::from_profile( + crate::profile::test_profile( + crate::RelayProfileKind::Public, + RelayUrlPolicy::Public, + ["wss://one.example", "wss://two.example"], + ) + .unwrap(), + ); + let transport = NostrTransport::with_client( + config, + Arc::new(MockRelayClient { + outcomes: BTreeMap::new(), + }), + ); + let request = request(); + let selected = TargetSet::new(vec![request.target_set().targets()[0].clone()]).unwrap(); + let prepared = transport + .prepare_delivery_selected(request.clone(), &selected) + .unwrap(); + assert_eq!(prepared.request(), &request); + assert_eq!(prepared.authorized.len(), 1); + assert_eq!(&prepared.authorized[0].1, &selected.targets()[0]); + let receipt = + futures::executor::block_on(transport.execute_prepared_delivery(prepared)).unwrap(); + receipt.validate_for_request(&request).unwrap(); + assert!(receipt.target_receipts()[0].was_attempted()); + assert_eq!( + receipt.target_receipts()[0].outcome().kind(), + DeliveryOutcomeKind::Accepted + ); + assert!(!receipt.target_receipts()[1].was_attempted()); + assert_eq!( + receipt.target_receipts()[1].outcome().code(), + Some("target_not_selected") + ); + assert!(receipt.target_receipts()[1].outcome().is_retryable()); + assert!(!receipt.is_satisfied(&request).unwrap()); + let foreign = + TargetSet::new(vec![Target::nostr_relay("wss://foreign.example").unwrap()]).unwrap(); + let failure = transport + .prepare_delivery_selected(request.clone(), &foreign) + .unwrap_err(); + failure.validate_for_request(&request).unwrap(); + assert_eq!(failure.code(), "invalid_transport_contract"); + let all = futures::executor::block_on( + transport.deliver_selected(request.clone(), request.target_set().clone()), + ) + .unwrap(); + assert!(all.is_satisfied(&request).unwrap()); + } + + #[test] fn upstream_messages_map_to_stable_outcomes() { let cases = [ ("duplicate: already have", DeliveryOutcomeKind::Accepted), diff --git a/crates/transport_nostr/tests/exact_delivery.rs b/crates/transport_nostr/tests/exact_delivery.rs @@ -90,6 +90,64 @@ async fn reply(socket: &mut WebSocketStream<TcpStream>, value: Value) { } #[tokio::test] +async fn selected_delivery_never_connects_to_held_target_and_preserves_full_request() { + let a = TcpListener::bind("127.0.0.1:0").await.unwrap(); + let b = TcpListener::bind("127.0.0.1:0").await.unwrap(); + let config = Config::from_profile( + RelayProfile::explicit( + RelayProfileKind::Simulator, + [a.local_addr().unwrap(), b.local_addr().unwrap()].map(|address| { + RelayEndpoint::new( + format!("ws://{address}"), + RelayUrlPolicy::Local, + RelayAccess::ReadWrite, + ) + .unwrap() + }), + ) + .unwrap(), + ) + .with_timeouts(1_000, 2_000, 500) + .unwrap(); + let raw = raw_event(); + let expected = format!("[\"EVENT\",{raw}]"); + let server = tokio::spawn(async move { + let (tcp, _) = a.accept().await.unwrap(); + let mut socket = accept_async(tcp).await.unwrap(); + let wire = text(&mut socket).await; + assert_eq!(wire, expected); + let event: Value = serde_json::from_str(&wire).unwrap(); + reply(&mut socket, json!(["OK", event[1]["id"], true, ""])).await; + }); + let transport = NostrTransport::new(config.clone()); + let request = request(&config, &raw); + let selected = TargetSet::new(vec![request.target_set().targets()[0].clone()]).unwrap(); + let result = transport + .deliver_selected(request.clone(), selected) + .await + .unwrap(); + result.validate_for_request(&request).unwrap(); + assert!(result.target_receipts()[0].was_attempted()); + assert!( + result.target_receipts()[0] + .outcome() + .satisfies(SatisfactionClass::Accepted) + ); + assert!(!result.target_receipts()[1].was_attempted()); + assert_eq!( + result.target_receipts()[1].outcome().code(), + Some("target_not_selected") + ); + assert!(!result.is_satisfied(&request).unwrap()); + server.await.unwrap(); + assert!( + tokio::time::timeout(Duration::from_millis(150), b.accept()) + .await + .is_err() + ); +} + +#[tokio::test] async fn exact_signed_bytes_read_requests_and_auth_use_the_same_connection() { let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); let config = config(&format!("ws://{}", listener.local_addr().unwrap()), 2_000);