selected.rs (5256B)
1 use super::*; 2 use radroots_transport::{ 3 DeliveryReceipt, DeliveryRequest, EventSink, SinkFailure, 4 outcome::DeliveryOutcome, 5 sink::{DeliveryTargetReceipt, SinkStatus}, 6 }; 7 8 #[derive(Default)] 9 struct SelectedSink(std::sync::Mutex<Vec<(DeliveryRequest, TargetSet)>>); 10 11 impl EventSink for SelectedSink { 12 fn status(&self) -> radroots_transport::BoxFuture<'_, Result<SinkStatus, TransportError>> { 13 Box::pin(async { panic!("selected delivery does not probe status") }) 14 } 15 16 fn deliver( 17 &self, 18 _: DeliveryRequest, 19 ) -> radroots_transport::BoxFuture<'_, Result<DeliveryReceipt, SinkFailure>> { 20 Box::pin(async { panic!("selected delivery must not widen to ordinary delivery") }) 21 } 22 23 fn deliver_selected( 24 &self, 25 request: DeliveryRequest, 26 selected: TargetSet, 27 ) -> radroots_transport::BoxFuture<'_, Result<DeliveryReceipt, SinkFailure>> { 28 Box::pin(async move { 29 request.validate_target_selection(&selected).unwrap(); 30 self.0 31 .lock() 32 .unwrap() 33 .push((request.clone(), selected.clone())); 34 Ok(DeliveryReceipt::for_request( 35 &request, 36 request 37 .target_set() 38 .targets() 39 .iter() 40 .map(|target| { 41 if selected.targets().contains(target) { 42 DeliveryTargetReceipt::attempted( 43 target.clone(), 44 DeliveryOutcome::accepted(), 45 ) 46 } else { 47 DeliveryTargetReceipt::skipped( 48 target.clone(), 49 DeliveryOutcome::unavailable(), 50 ) 51 .unwrap() 52 } 53 }) 54 .collect(), 55 ) 56 .unwrap()) 57 }) 58 } 59 } 60 61 #[tokio::test] 62 async fn sdk_selected_delivery_preserves_subset_and_full_durable_request() { 63 let storage = Arc::new(MemoryStorage::new( 64 SourceGeneration::new([211; 32]).unwrap(), 65 )); 66 let signer = radroots_nostr::signing::LocalSigner::new( 67 radroots_nostr::key::SecretKey::parse( 68 "0000000000000000000000000000000000000000000000000000000000000001", 69 ) 70 .unwrap(), 71 ) 72 .unwrap(); 73 let author = "79be667ef9dcbbac55a06295ce870b07029bfcdb2dce28d959f2815b16f81798"; 74 let (clock, ids, deadlines) = HostPolicy::default().composition(); 75 let now = clock.now_unix_ms().unwrap(); 76 let sink = Arc::new(SelectedSink::default()); 77 let engine = Engine::builder(storage.clone(), clock, ids, deadlines) 78 .signer(Arc::new(signer)) 79 .sink(sink.clone()) 80 .build() 81 .unwrap(); 82 let client = ClientBuilder::new() 83 .storage(storage) 84 .sync_engine(engine) 85 .build() 86 .unwrap(); 87 let operations = client.sync().unwrap().unwrap(); 88 let targets = TargetSet::new(vec![ 89 target(), 90 Target::nostr_relay("wss://held.example").unwrap(), 91 ]) 92 .unwrap(); 93 let selected = TargetSet::new(vec![targets.targets()[0].clone()]).unwrap(); 94 let request = PushRequest::new( 95 SyncId::new([211; 16]).unwrap(), 96 IdempotencyKey::parse("sdk-selected-delivery").unwrap(), 97 Actor::new( 98 PublicKey::from_hex(author).unwrap(), 99 ActorSource::ExplicitPublicKey, 100 [AuthorRole::Any], 101 ) 102 .unwrap(), 103 AuthoredEventPlan::from_generic( 104 GenericEventDraft::new( 105 "radroots.social.geochat.v1", 106 20_000, 107 now / 1_000, 108 Vec::new(), 109 "selected", 110 author, 111 ) 112 .unwrap(), 113 ) 114 .unwrap(), 115 targets.clone(), 116 SatisfactionPolicy::new(SatisfactionClass::Accepted, TargetPolicy::all()), 117 now + 60_000, 118 CancellationPolicy::LocalCooperative, 119 ) 120 .unwrap(); 121 let id = request.operation_id(); 122 operations.submit_push(request).await.unwrap(); 123 let original = operations.push_status(id).await.unwrap().unwrap(); 124 let original_request = original.delivery_plan().request().unwrap(); 125 assert_eq!( 126 operations 127 .deliver_push_selected( 128 id, 129 TargetSet::new(vec![Target::nostr_relay("wss://foreign.example").unwrap()]) 130 .unwrap(), 131 ) 132 .await 133 .unwrap_err(), 134 Error::InvalidDeliveryRequest 135 ); 136 assert!(sink.0.lock().unwrap().is_empty()); 137 operations 138 .deliver_push_selected(id, selected.clone()) 139 .await 140 .unwrap(); 141 let after = operations.push_status(id).await.unwrap().unwrap(); 142 assert_eq!(after.delivery_plan().request(), Some(original_request)); 143 assert_eq!( 144 sink.0.lock().unwrap().as_slice(), 145 &[(original_request.clone(), selected)] 146 ); 147 assert!(!after.delivery_plan().state().is_terminal()); 148 assert_eq!(after.delivery_plan().delivery_facts().len(), 1); 149 client.close().await.unwrap(); 150 }