1use std::collections::HashMap;
14use std::collections::VecDeque;
15use std::ops::Deref;
16use std::sync::Arc;
17use std::sync::Mutex;
18use std::sync::RwLock;
19
20use bytes::Bytes;
21use futures::lock::Mutex as AsyncMutex;
22use rings_core::dht::Did;
23
24use super::Ctx;
25use super::Envelope;
26use super::Inbound;
27use super::Interpret;
28use super::MaybeSend;
29use super::Protocol;
30use super::Reject;
31use super::Transition;
32use super::Wire;
33use crate::error::Error;
34use crate::error::Result;
35use crate::processor::Processor;
36use crate::sync_lock::lock;
37
38const MAX_FIXPOINT_STEPS: u32 = 1024;
41
42#[cfg(rings_native)]
44pub(crate) type DynHandler = dyn Handler + Send + Sync;
45#[cfg(rings_browser)]
47pub(crate) type DynHandler = dyn Handler;
48
49type HandlerMap = RwLock<HashMap<String, Arc<DynHandler>>>;
50
51#[cfg_attr(rings_browser, async_trait::async_trait(?Send))]
54#[cfg_attr(rings_native, async_trait::async_trait)]
55pub(crate) trait Handler {
56 async fn handle(&self, core: &Core, from: Did, payload: Bytes) -> Result<Vec<Inbound>>;
59}
60
61#[derive(Clone)]
65pub(crate) struct Core {
66 processor: Arc<Processor>,
67 handlers: Arc<HandlerMap>,
68}
69
70impl Core {
71 pub fn did(&self) -> Did {
73 self.processor.did()
74 }
75
76 pub async fn send(&self, to: Did, namespace: &str, payload: Bytes) -> Result<()> {
78 let envelope = Envelope::new(namespace, payload);
79 self.processor.send_envelope(to, &envelope).await?;
80 Ok(())
81 }
82
83 async fn send_direct(&self, to: Did, namespace: &str, payload: Bytes) -> Result<()> {
85 let envelope = Envelope::new(namespace, payload);
86 self.processor.send_direct_envelope(to, &envelope).await?;
87 Ok(())
88 }
89
90 pub async fn inject(&self, namespace: &str, payload: Bytes) -> Result<()> {
93 self.dispatch(self.did(), Envelope::new(namespace, payload))
94 .await
95 }
96
97 pub(crate) async fn dispatch(&self, from: Did, envelope: Envelope) -> Result<()> {
108 let mut queue: VecDeque<Inbound> = VecDeque::new();
109 queue.push_back(Inbound {
110 namespace: envelope.namespace,
111 from,
112 payload: envelope.payload,
113 });
114
115 let mut budget = MAX_FIXPOINT_STEPS;
116 while let Some(Inbound {
117 namespace,
118 from,
119 payload,
120 }) = queue.pop_front()
121 {
122 if budget == 0 {
123 return Err(Error::ExtensionError(format!(
124 "fixpoint budget ({MAX_FIXPOINT_STEPS}) exhausted; last namespace {namespace:?}"
125 )));
126 }
127 budget -= 1;
128
129 match self.handler(namespace.as_str()) {
130 Some(handler) => queue.extend(handler.handle(self, from, payload).await?),
131 None => tracing::debug!(
132 "no protocol registered for namespace {:?}, dropping",
133 namespace
134 ),
135 }
136 }
137 Ok(())
138 }
139
140 fn handler(&self, namespace: &str) -> Option<Arc<DynHandler>> {
141 self.handlers.read().ok()?.get(namespace).map(Arc::clone)
142 }
143}
144
145#[derive(Clone)]
153pub struct Scope {
154 core: Core,
155 namespace: String,
156}
157
158impl Scope {
159 pub(crate) fn new(core: Core, namespace: String) -> Self {
161 Self { core, namespace }
162 }
163
164 pub fn did(&self) -> Did {
166 self.core.did()
167 }
168
169 pub fn namespace(&self) -> &str {
171 self.namespace.as_str()
172 }
173
174 pub async fn send(&self, to: Did, payload: Bytes) -> Result<()> {
176 self.core.send(to, self.namespace.as_str(), payload).await
177 }
178
179 pub(crate) async fn send_direct(&self, to: Did, payload: Bytes) -> Result<()> {
184 self.core
185 .send_direct(to, self.namespace.as_str(), payload)
186 .await
187 }
188
189 pub(crate) async fn inject(&self, payload: Bytes) -> Result<()> {
201 self.core.inject(self.namespace.as_str(), payload).await
202 }
203}
204
205pub struct EffectScope {
211 scope: Scope,
212}
213
214impl EffectScope {
215 pub(crate) fn new(scope: Scope) -> Self {
216 Self { scope }
217 }
218
219 pub fn did(&self) -> Did {
221 self.scope.did()
222 }
223
224 pub fn namespace(&self) -> &str {
226 self.scope.namespace()
227 }
228
229 pub async fn send(&self, to: Did, payload: Bytes) -> Result<()> {
231 self.scope.send(to, payload).await
232 }
233
234 pub(crate) fn lifecycle(&self) -> Scope {
239 self.scope.clone()
240 }
241}
242
243struct Runner<P: Protocol, I> {
246 protocol: P,
247 interpret: I,
248 state: Mutex<P::State>,
249 transition_gate: AsyncMutex<()>,
250 #[cfg(all(test, rings_native))]
251 after_decode_for_test: Option<Arc<dyn Fn() + Send + Sync>>,
252 #[cfg(all(test, rings_native))]
253 after_commit_for_test: Option<Arc<dyn Fn() + Send + Sync>>,
254 #[cfg(all(test, rings_native))]
255 before_gate_wait_for_test: Option<Arc<dyn Fn(bool) + Send + Sync>>,
256}
257
258#[cfg_attr(rings_browser, async_trait::async_trait(?Send))]
259#[cfg_attr(rings_native, async_trait::async_trait)]
260impl<P, I> Handler for Runner<P, I>
261where
262 P: Protocol + MaybeSend + 'static,
263 P::State: MaybeSend + 'static,
264 P::Effect: MaybeSend,
265 I: Interpret<Effect = P::Effect> + MaybeSend + 'static,
266{
267 async fn handle(&self, core: &Core, from: Did, payload: Bytes) -> Result<Vec<Inbound>> {
268 let event = match self.protocol.decode(Wire {
271 from,
272 me: core.did(),
273 payload: payload.as_ref(),
274 }) {
275 Ok(event) => event,
276 Err(Reject(why)) => {
277 tracing::debug!("drop on {}: {why}", self.protocol.namespace());
278 return Ok(Vec::new());
279 }
280 };
281
282 #[cfg(all(test, rings_native))]
283 if let Some(observe) = self.after_decode_for_test.as_ref() {
284 observe();
285 }
286
287 #[cfg(all(test, rings_native))]
293 if let Some(observe) = self.before_gate_wait_for_test.as_ref() {
294 observe(self.transition_gate.try_lock().is_none());
297 }
298 let _transition_turn = self.transition_gate.lock().await;
299
300 let namespace = self.protocol.namespace().to_string();
305 let scope = EffectScope::new(Scope::new(core.clone(), namespace.clone()));
306 let mut feedback = VecDeque::new();
307 feedback.push_back(event);
308 let mut feedback_budget = MAX_FIXPOINT_STEPS;
309 while let Some(event) = feedback.pop_front() {
310 if feedback_budget == 0 {
311 return Err(Error::ExtensionError(format!(
312 "feedback fixpoint budget ({MAX_FIXPOINT_STEPS}) exhausted on {namespace:?}"
313 )));
314 }
315 feedback_budget -= 1;
316
317 let effects = {
321 let mut guard = lock(&self.state)?;
322 let Transition { state, effects } = self.protocol.step(
323 Ctx {
324 did: core.did(),
325 state: guard.deref(),
326 },
327 event,
328 );
329 *guard = state;
330 effects
331 };
332
333 #[cfg(all(test, rings_native))]
334 if let Some(observe) = self.after_commit_for_test.as_ref() {
335 observe();
336 }
337
338 for effect in effects {
339 for payload in self.interpret.run(&scope, effect).await? {
340 match self.protocol.decode(Wire {
341 from: core.did(),
342 me: core.did(),
343 payload: payload.as_ref(),
344 }) {
345 Ok(event) => feedback.push_back(event),
346 Err(Reject(why)) => {
347 tracing::debug!("drop feedback on {}: {why}", self.protocol.namespace())
348 }
349 }
350 }
351 }
352 }
353 Ok(Vec::new())
354 }
355}
356
357#[derive(Clone)]
361pub struct Extensions {
362 core: Core,
363}
364
365impl Extensions {
366 pub fn new(processor: Arc<Processor>) -> Self {
368 Self {
369 core: Core {
370 processor,
371 handlers: Arc::new(RwLock::new(HashMap::new())),
372 },
373 }
374 }
375
376 pub(crate) fn core(&self) -> Core {
383 self.core.clone()
384 }
385
386 pub fn register<P, I>(&self, protocol: P, interpret: I) -> Result<()>
390 where
391 P: Protocol + MaybeSend + 'static,
392 P::State: MaybeSend + 'static,
393 P::Effect: MaybeSend,
394 I: Interpret<Effect = P::Effect> + MaybeSend + 'static,
395 {
396 self.insert(protocol, interpret, false)
397 }
398
399 pub fn replace<P, I>(&self, protocol: P, interpret: I) -> Result<()>
402 where
403 P: Protocol + MaybeSend + 'static,
404 P::State: MaybeSend + 'static,
405 P::Effect: MaybeSend,
406 I: Interpret<Effect = P::Effect> + MaybeSend + 'static,
407 {
408 self.insert(protocol, interpret, true)
409 }
410
411 pub fn register_many<P, I>(&self, items: Vec<(P, I)>) -> Result<()>
417 where
418 P: Protocol + MaybeSend + 'static,
419 P::State: MaybeSend + 'static,
420 P::Effect: MaybeSend,
421 I: Interpret<Effect = P::Effect> + MaybeSend + 'static,
422 {
423 let prepared: Vec<(String, Arc<DynHandler>, Vec<&'static str>)> = items
425 .into_iter()
426 .map(|(protocol, interpret)| {
427 let capabilities = protocol.capabilities().to_vec();
428 let namespace = protocol.namespace().to_string();
429 let state = Mutex::new(protocol.init());
430 let runner: Arc<DynHandler> = Arc::new(Runner {
431 protocol,
432 interpret,
433 state,
434 transition_gate: AsyncMutex::new(()),
435 #[cfg(all(test, rings_native))]
436 after_decode_for_test: None,
437 #[cfg(all(test, rings_native))]
438 after_commit_for_test: None,
439 #[cfg(all(test, rings_native))]
440 before_gate_wait_for_test: None,
441 });
442 (namespace, runner, capabilities)
443 })
444 .collect();
445
446 let mut handlers = self.core.handlers.write().map_err(|_| Error::Lock)?;
447 for (index, (namespace, _, _)) in prepared.iter().enumerate() {
449 let duplicate_in_batch = prepared
450 .iter()
451 .take(index)
452 .any(|(seen, _, _)| seen == namespace);
453 if duplicate_in_batch || handlers.contains_key(namespace) {
454 return Err(Error::ExtensionError(format!(
455 "namespace {namespace:?} is already registered"
456 )));
457 }
458 }
459 self.core.processor.add_online_node_capabilities(
460 prepared
461 .iter()
462 .flat_map(|(_, _, capabilities)| capabilities.iter().copied()),
463 )?;
464 for (namespace, runner, _) in prepared {
466 handlers.insert(namespace, runner);
467 }
468 Ok(())
469 }
470
471 fn insert<P, I>(&self, protocol: P, interpret: I, replace: bool) -> Result<()>
472 where
473 P: Protocol + MaybeSend + 'static,
474 P::State: MaybeSend + 'static,
475 P::Effect: MaybeSend,
476 I: Interpret<Effect = P::Effect> + MaybeSend + 'static,
477 {
478 let capabilities = protocol.capabilities();
479 let namespace = protocol.namespace().to_string();
480 let state = Mutex::new(protocol.init());
481 let runner: Arc<DynHandler> = Arc::new(Runner {
482 protocol,
483 interpret,
484 state,
485 transition_gate: AsyncMutex::new(()),
486 #[cfg(all(test, rings_native))]
487 after_decode_for_test: None,
488 #[cfg(all(test, rings_native))]
489 after_commit_for_test: None,
490 #[cfg(all(test, rings_native))]
491 before_gate_wait_for_test: None,
492 });
493 let mut handlers = self.core.handlers.write().map_err(|_| Error::Lock)?;
494 if !replace && handlers.contains_key(&namespace) {
495 return Err(Error::ExtensionError(format!(
496 "namespace {namespace:?} is already registered"
497 )));
498 }
499 self.core
500 .processor
501 .add_online_node_capabilities(capabilities.iter().copied())?;
502 handlers.insert(namespace, runner);
503 Ok(())
504 }
505
506 pub fn contains(&self, namespace: &str) -> bool {
508 self.core
509 .handlers
510 .read()
511 .map(|h| h.contains_key(namespace))
512 .unwrap_or(false)
513 }
514
515 pub(crate) async fn dispatch(&self, from: Did, envelope: Envelope) -> Result<()> {
521 self.core.dispatch(from, envelope).await
522 }
523}
524
525#[cfg(all(test, rings_native))]
526mod tests {
527 use std::collections::HashMap;
528 use std::net::SocketAddr;
529 use std::sync::Arc;
530 use std::sync::Mutex;
531
532 use async_trait::async_trait;
533 use rings_core::ecc::SecretKey;
534 use rings_core::session::SessionSk;
535 use tokio::sync::Notify;
536
537 use super::*;
538 use crate::extension::protocols::relay::ControlSendTestHook;
539 use crate::extension::protocols::relay::NativeRelay;
540 use crate::extension::protocols::relay::Relay;
541 use crate::extension::protocols::relay::RelayCommand;
542 use crate::extension::protocols::relay::RelayEffect;
543 use crate::extension::protocols::relay::TCP;
544 use crate::extension::transport::engine::TransportSessions;
545 use crate::extension::transport::Frame;
546 use crate::extension::transport::Initiator;
547 use crate::extension::transport::SessionId;
548 use crate::extension::transport::SessionKey;
549 use crate::processor::ProcessorBuilder;
550 use crate::processor::ProcessorConfig;
551
552 struct OrderedProtocol;
553
554 impl Protocol for OrderedProtocol {
555 type State = u8;
556 type Event = u8;
557 type Effect = u8;
558
559 fn namespace(&self) -> &str {
560 "ordered-effects"
561 }
562
563 fn init(&self) -> Self::State {
564 0
565 }
566
567 fn decode(&self, wire: Wire<'_>) -> std::result::Result<Self::Event, Reject> {
568 let event = wire
569 .payload
570 .first()
571 .copied()
572 .ok_or_else(|| Reject("missing effect value".to_string()))?;
573 Ok(event)
574 }
575
576 fn step(
577 &self,
578 ctx: Ctx<'_, Self::State>,
579 event: Self::Event,
580 ) -> Transition<Self::State, Self::Effect> {
581 Transition::with(ctx.state.saturating_add(1), vec![event])
582 }
583 }
584
585 #[derive(Default)]
586 struct BlockingOrderedInterpreter {
587 first_effect_started: Notify,
588 release_first_effect: Notify,
589 observed: Mutex<Vec<u8>>,
590 }
591
592 #[async_trait]
593 impl Interpret for Arc<BlockingOrderedInterpreter> {
594 type Effect = u8;
595
596 async fn run(&self, _scope: &EffectScope, effect: Self::Effect) -> Result<Vec<Bytes>> {
597 if effect == 1 {
598 self.first_effect_started.notify_one();
599 self.release_first_effect.notified().await;
600 }
601 lock(&self.observed)?.push(effect);
602 Ok(Vec::new())
603 }
604 }
605
606 #[derive(Default)]
607 struct RelayFeedbackInterpreter {
608 first_effect_started: Notify,
609 release_first_effect: Notify,
610 first_connect_seen: Mutex<bool>,
611 observed_connects: Mutex<Vec<SessionId>>,
612 }
613
614 #[async_trait]
615 impl Interpret for Arc<RelayFeedbackInterpreter> {
616 type Effect = RelayEffect<SocketAddr>;
617
618 async fn run(&self, _scope: &EffectScope, effect: Self::Effect) -> Result<Vec<Bytes>> {
619 match effect {
620 RelayEffect::Connect { key, .. } => {
621 let first_connect = {
622 let mut seen = lock(&self.first_connect_seen)?;
623 let first_connect = !*seen;
624 *seen = true;
625 first_connect
626 };
627 if first_connect {
628 self.first_effect_started.notify_one();
629 self.release_first_effect.notified().await;
630 let feedback = RelayCommand::<SocketAddr>::Untrack {
631 peer: key.peer,
632 session: key.session,
633 initiator: key.initiator,
634 };
635 return rings_codec::serialize(&feedback)
636 .map(Bytes::from)
637 .map(|payload| vec![payload])
638 .map_err(|_| Error::EncodeError);
639 }
640 lock(&self.observed_connects)?.push(key.session);
641 Ok(Vec::new())
642 }
643 _ => Ok(Vec::new()),
644 }
645 }
646 }
647
648 #[derive(Default)]
649 struct FailingOrderedInterpreter {
650 observed: Mutex<Vec<u8>>,
651 }
652
653 #[async_trait]
654 impl Interpret for Arc<FailingOrderedInterpreter> {
655 type Effect = u8;
656
657 async fn run(&self, _scope: &EffectScope, effect: Self::Effect) -> Result<Vec<Bytes>> {
658 lock(&self.observed)?.push(effect);
659 if effect == 1 {
660 return Err(Error::ExtensionError(
661 "intentional effect failure".to_string(),
662 ));
663 }
664 Ok(Vec::new())
665 }
666 }
667
668 fn extensions() -> Result<Extensions> {
669 let session = SessionSk::new_with_seckey(&SecretKey::random())?;
670 let config = ProcessorConfig::new(1, String::new(), session, 1);
671 let processor = ProcessorBuilder::from_config(&config)?
672 .advertise_presence(false)
673 .build()?;
674 Ok(Extensions::new(Arc::new(processor)))
675 }
676
677 #[tokio::test]
678 async fn test_unknown_legacy_namespace_is_a_nonfatal_drop() -> Result<()> {
679 let extensions = extensions()?;
680 let from = extensions.core().did();
681
682 extensions
683 .dispatch(
684 from,
685 Envelope::new("snark", Bytes::from_static(b"legacy-task")),
686 )
687 .await?;
688
689 assert!(extensions.core.handler("snark").is_none());
690 Ok(())
691 }
692
693 #[tokio::test]
694 async fn test_committed_transitions_execute_effects_in_commit_order() -> Result<()> {
695 let extensions = extensions()?;
698 let interpreter = Arc::new(BlockingOrderedInterpreter::default());
699 let gate_wait = Arc::new(Notify::new());
700 let gate_contention = Arc::new(Mutex::new(Vec::new()));
701 let gate_observer = {
702 let gate_wait = Arc::clone(&gate_wait);
703 let gate_contention = Arc::clone(&gate_contention);
704 Arc::new(move |contended| {
705 gate_contention
706 .lock()
707 .expect("test gate witness lock")
708 .push(contended);
709 gate_wait.notify_one();
710 }) as Arc<dyn Fn(bool) + Send + Sync>
711 };
712 let committed = Arc::new(Mutex::new(0_u8));
713 let commit_observer = {
714 let committed = Arc::clone(&committed);
715 Arc::new(move || {
716 *committed.lock().expect("test commit witness lock") += 1;
717 }) as Arc<dyn Fn() + Send + Sync>
718 };
719 let runner: Arc<DynHandler> = Arc::new(Runner {
720 protocol: OrderedProtocol,
721 interpret: Arc::clone(&interpreter),
722 state: Mutex::new(0),
723 transition_gate: AsyncMutex::new(()),
724 after_decode_for_test: None,
725 after_commit_for_test: Some(commit_observer),
726 before_gate_wait_for_test: Some(gate_observer),
727 });
728 extensions
729 .core
730 .handlers
731 .write()
732 .map_err(|_| Error::Lock)?
733 .insert("ordered-effects".to_string(), runner);
734 let from = extensions.core().did();
735
736 let first_extensions = extensions.clone();
737 let first = tokio::spawn(async move {
738 first_extensions
739 .dispatch(
740 from,
741 Envelope::new("ordered-effects", Bytes::from_static(&[1])),
742 )
743 .await
744 });
745 interpreter.first_effect_started.notified().await;
746 gate_wait.notified().await;
747
748 let second_extensions = extensions.clone();
749 let second = tokio::spawn(async move {
750 second_extensions
751 .dispatch(
752 from,
753 Envelope::new("ordered-effects", Bytes::from_static(&[2])),
754 )
755 .await
756 });
757 gate_wait.notified().await;
758 assert_eq!(*lock(&gate_contention)?, vec![false, true]);
759 assert!(!second.is_finished());
760 assert!(lock(&interpreter.observed)?.is_empty());
761 assert_eq!(*lock(&committed)?, 1);
762
763 interpreter.release_first_effect.notify_one();
764 first
765 .await
766 .map_err(|error| Error::ExtensionError(error.to_string()))??;
767 second
768 .await
769 .map_err(|error| Error::ExtensionError(error.to_string()))??;
770 assert_eq!(*lock(&committed)?, 2);
771 assert_eq!(*lock(&interpreter.observed)?, vec![1, 2]);
772 Ok(())
773 }
774
775 #[tokio::test]
776 async fn test_failed_effect_releases_ordered_turn_for_later_transition() -> Result<()> {
777 let extensions = extensions()?;
780 let interpreter = Arc::new(FailingOrderedInterpreter::default());
781 extensions.register(OrderedProtocol, Arc::clone(&interpreter))?;
782 let from = extensions.core().did();
783
784 let failed = extensions
785 .dispatch(
786 from,
787 Envelope::new("ordered-effects", Bytes::from_static(&[1])),
788 )
789 .await;
790 assert!(matches!(failed, Err(Error::ExtensionError(_))));
791 extensions
792 .dispatch(
793 from,
794 Envelope::new("ordered-effects", Bytes::from_static(&[2])),
795 )
796 .await?;
797
798 assert_eq!(*lock(&interpreter.observed)?, vec![1, 2]);
799 Ok(())
800 }
801
802 #[tokio::test]
803 async fn test_returned_feedback_precedes_a_waiting_transition() -> Result<()> {
804 let extensions = extensions()?;
808 let interpreter = Arc::new(RelayFeedbackInterpreter::default());
809 let decoded = Arc::new(Notify::new());
810 let observer = {
811 let decoded = Arc::clone(&decoded);
812 Arc::new(move || decoded.notify_one()) as Arc<dyn Fn() + Send + Sync>
813 };
814 let protocol = Relay::tcp(HashMap::from([(
815 "web".to_string(),
816 "127.0.0.1:80"
817 .parse::<SocketAddr>()
818 .map_err(|error| Error::ExtensionError(error.to_string()))?,
819 )]));
820 let state = protocol.init();
821 let runner: Arc<DynHandler> = Arc::new(Runner {
822 protocol,
823 interpret: Arc::clone(&interpreter),
824 state: Mutex::new(state),
825 transition_gate: AsyncMutex::new(()),
826 after_decode_for_test: Some(observer),
827 after_commit_for_test: None,
828 before_gate_wait_for_test: None,
829 });
830 extensions
831 .core
832 .handlers
833 .write()
834 .map_err(|_| Error::Lock)?
835 .insert(TCP.to_string(), runner);
836 let from: Did = SecretKey::random().address().into();
837 let open = rings_codec::serialize(&Frame::Open {
838 session: SessionId(0),
839 service: "web".to_string(),
840 })
841 .map(Bytes::from)
842 .map_err(|_| Error::EncodeError)?;
843 let first_open = open.clone();
844
845 let first_extensions = extensions.clone();
846 let first = tokio::spawn(async move {
847 first_extensions
848 .dispatch(from, Envelope::new(TCP, first_open))
849 .await
850 });
851 interpreter.first_effect_started.notified().await;
852 decoded.notified().await;
853
854 let second_extensions = extensions.clone();
855 let second = tokio::spawn(async move {
856 second_extensions
857 .dispatch(from, Envelope::new(TCP, open))
858 .await
859 });
860 decoded.notified().await;
861 assert!(!second.is_finished());
862
863 interpreter.release_first_effect.notify_one();
864 let timeout = std::time::Duration::from_secs(1);
865 tokio::time::timeout(timeout, first)
866 .await
867 .map_err(|_| Error::ExtensionError("first feedback turn timed out".to_string()))?
868 .map_err(|error| Error::ExtensionError(error.to_string()))??;
869 tokio::time::timeout(timeout, second)
870 .await
871 .map_err(|_| Error::ExtensionError("second feedback turn timed out".to_string()))?
872 .map_err(|error| Error::ExtensionError(error.to_string()))??;
873 assert_eq!(*lock(&interpreter.observed_connects)?, vec![SessionId(0)]);
874 Ok(())
875 }
876
877 #[tokio::test]
878 async fn test_missing_open_accepted_resource_returns_synchronous_untrack() -> Result<()> {
879 let extensions = extensions()?;
880 let effect_scope = EffectScope::new(Scope::new(extensions.core(), TCP.to_string()));
881 let interpreter = NativeRelay::new(Arc::new(TransportSessions::new()));
882 let peer: Did = SecretKey::random().address().into();
883 let key = SessionKey::new(peer, TCP, SessionId(9), Initiator::Local);
884
885 let feedback = interpreter
886 .run(&effect_scope, RelayEffect::OpenAccepted {
887 token: 77,
888 key: key.clone(),
889 service: "missing-pending-resource".to_string(),
890 })
891 .await?;
892
893 assert_eq!(feedback.len(), 1);
894 assert!(matches!(
895 rings_codec::deserialize::<RelayCommand<SocketAddr>>(feedback[0].as_ref()),
896 Ok(RelayCommand::Untrack {
897 peer: actual_peer,
898 session: SessionId(9),
899 initiator: Initiator::Local,
900 }) if actual_peer == peer
901 ));
902 Ok(())
903 }
904
905 #[tokio::test]
906 async fn test_terminal_relay_control_effect_does_not_await_overlay_send() -> Result<()> {
907 let extensions = extensions()?;
908 let hook = Arc::new(ControlSendTestHook::default());
909 let interpreter = Arc::new(NativeRelay::new_with_control_send_test_hook(
910 Arc::new(TransportSessions::new()),
911 Arc::clone(&hook),
912 ));
913 let peer: Did = SecretKey::random().address().into();
914 let core = extensions.core();
915 let application = tokio::spawn(async move {
916 let effect_scope = EffectScope::new(Scope::new(core, TCP.to_string()));
917 interpreter
918 .run(&effect_scope, RelayEffect::SendClose {
919 to: peer,
920 session: SessionId(5),
921 from_opener: false,
922 })
923 .await
924 });
925
926 tokio::time::timeout(std::time::Duration::from_secs(1), hook.wait_until_blocked())
930 .await
931 .map_err(|_| {
932 Error::ExtensionError("control outbox did not reach test gate".to_string())
933 })?;
934 let applied = tokio::time::timeout(std::time::Duration::from_secs(1), application)
935 .await
936 .map_err(|_| {
937 Error::ExtensionError("terminal control effect held the gate".to_string())
938 })?
939 .map_err(|error| Error::ExtensionError(error.to_string()))??;
940
941 assert!(applied.is_empty());
942 hook.release();
943 Ok(())
944 }
945
946 #[tokio::test]
947 async fn test_saturated_peer_control_lane_does_not_block_another_peer() -> Result<()> {
948 let extensions = extensions()?;
949 let hook = Arc::new(ControlSendTestHook::default());
950 let interpreter = NativeRelay::new_with_control_send_test_hook(
951 Arc::new(TransportSessions::new()),
952 Arc::clone(&hook),
953 );
954 let blocked_peer: Did = SecretKey::random().address().into();
955 let independent_peer: Did = SecretKey::random().address().into();
956 let effect_scope = EffectScope::new(Scope::new(extensions.core(), TCP.to_string()));
957
958 interpreter
959 .run(&effect_scope, RelayEffect::SendClose {
960 to: blocked_peer,
961 session: SessionId(0),
962 from_opener: false,
963 })
964 .await?;
965 tokio::time::timeout(std::time::Duration::from_secs(1), hook.wait_until_blocked())
966 .await
967 .map_err(|_| {
968 Error::ExtensionError("first peer control lane did not block".to_string())
969 })?;
970
971 let mut saturated = false;
972 for session in 1..=8 {
973 let result = interpreter
974 .run(&effect_scope, RelayEffect::SendClose {
975 to: blocked_peer,
976 session: SessionId(session),
977 from_opener: false,
978 })
979 .await;
980 if result.is_err() {
981 saturated = true;
982 break;
983 }
984 }
985 assert!(saturated, "the blocked peer must have a finite lane budget");
986
987 interpreter
988 .run(&effect_scope, RelayEffect::SendClose {
989 to: independent_peer,
990 session: SessionId(9),
991 from_opener: false,
992 })
993 .await?;
994 tokio::time::timeout(
995 std::time::Duration::from_secs(1),
996 hook.wait_until_completed(independent_peer),
997 )
998 .await
999 .map_err(|_| {
1000 Error::ExtensionError("independent peer control lane was blocked".to_string())
1001 })??;
1002
1003 hook.release();
1004 tokio::time::timeout(
1005 std::time::Duration::from_secs(1),
1006 hook.wait_until_completed(blocked_peer),
1007 )
1008 .await
1009 .map_err(|_| Error::ExtensionError("blocked peer lane did not resume".to_string()))??;
1010 Ok(())
1011 }
1012}