1use std::collections::HashMap;
37use std::collections::HashSet;
38use std::sync::Arc;
39
40use bytes::Bytes;
41use rings_core::dht::Did;
42use serde::de::DeserializeOwned;
43use serde::Deserialize;
44use serde::Serialize;
45
46use crate::extension::ext::Ctx;
47#[cfg(any(rings_native, rings_browser))]
48use crate::extension::ext::EffectScope;
49#[cfg(any(rings_native, rings_browser))]
50use crate::extension::ext::Interpret;
51use crate::extension::ext::MaybeSend;
52use crate::extension::ext::Protocol;
53use crate::extension::ext::Reject;
54use crate::extension::ext::Scope;
55use crate::extension::ext::Transition;
56use crate::extension::ext::Wire;
57use crate::extension::transport::EffectEnqueue;
58use crate::extension::transport::Frame;
59use crate::extension::transport::Initiator;
60use crate::extension::transport::SessionId;
61use crate::extension::transport::SessionKey;
62use crate::extension::transport::TransportKind;
63use crate::peer_quota::PeerQuota;
64
65#[cfg(any(rings_native, rings_browser))]
66mod control_outbox;
67#[cfg(any(rings_native, rings_browser))]
68use self::control_outbox::ControlOutbox;
69#[cfg(all(test, rings_native))]
70pub(crate) use self::control_outbox::ControlSendTestHook;
71
72pub const TCP: &str = "tcp";
74pub const UDP: &str = "udp";
76
77pub(crate) const MAX_RELAY_SESSIONS: usize = 1_024;
80pub(crate) const MAX_RELAY_SESSIONS_PER_PEER: usize = 64;
82#[derive(Clone, Debug, Serialize, Deserialize, PartialEq, Eq)]
85pub enum RelayCommand<T> {
86 RegisterService {
88 name: String,
90 target: T,
92 },
93 Accepted {
99 token: u64,
101 peer: Did,
103 service: String,
105 },
106 Untrack {
109 peer: Did,
111 session: SessionId,
113 initiator: Initiator,
115 },
116 Abort {
119 peer: Did,
121 session: SessionId,
123 initiator: Initiator,
125 },
126}
127
128pub enum RelayEvent<T> {
131 Command(RelayCommand<T>),
133 Frame {
135 from: Did,
137 frame: Frame,
139 },
140}
141
142#[derive(Clone, Debug, PartialEq, Eq)]
144pub enum RelayEffect<T> {
145 Connect {
147 key: SessionKey,
149 target: T,
151 kind: TransportKind,
153 },
154 Write {
156 key: SessionKey,
158 bytes: Bytes,
160 },
161 Shutdown {
163 key: SessionKey,
165 },
166 Close {
168 key: SessionKey,
170 },
171 SendClose {
175 to: Did,
177 session: SessionId,
179 from_opener: bool,
181 },
182 OpenAccepted {
186 token: u64,
188 key: SessionKey,
190 service: String,
192 },
193 RejectAccepted {
195 token: u64,
197 },
198}
199
200#[derive(Clone)]
204pub struct RelayState<T> {
205 services: Arc<HashMap<String, T>>,
206 sessions: Arc<HashSet<SessionKey>>,
209 session_quota: Arc<PeerQuota>,
215 peer_shutdown: Arc<HashSet<SessionKey>>,
220 next_session: u64,
224}
225
226impl<T> Default for RelayState<T> {
227 fn default() -> Self {
228 Self {
229 services: Arc::new(HashMap::new()),
230 sessions: Arc::new(HashSet::new()),
231 session_quota: Arc::new(PeerQuota::new(
232 MAX_RELAY_SESSIONS,
233 MAX_RELAY_SESSIONS_PER_PEER,
234 )),
235 peer_shutdown: Arc::new(HashSet::new()),
236 next_session: 0,
237 }
238 }
239}
240
241impl<T> RelayState<T> {
242 fn can_admit_session(&self, key: &SessionKey) -> bool {
243 self.session_quota.can_reserve(key.peer).is_ok()
244 }
245
246 fn insert_session(&mut self, key: SessionKey) -> bool {
247 if self.sessions.contains(&key) {
248 return false;
249 }
250 if Arc::make_mut(&mut self.session_quota)
251 .reserve(key.peer)
252 .is_err()
253 {
254 return false;
255 }
256 if Arc::make_mut(&mut self.sessions).insert(key.clone()) {
257 Arc::make_mut(&mut self.peer_shutdown).remove(&key);
258 true
259 } else {
260 let rolled_back = Arc::make_mut(&mut self.session_quota).release(key.peer);
261 debug_assert!(rolled_back);
262 false
263 }
264 }
265
266 fn remove_session(&mut self, key: &SessionKey) -> bool {
267 if !self.sessions.contains(key) {
268 return false;
269 }
270 if !Arc::make_mut(&mut self.session_quota).release(key.peer) {
271 debug_assert!(false, "session quota missing admitted peer {}", key.peer);
272 return false;
273 }
274 Arc::make_mut(&mut self.peer_shutdown).remove(key);
275 let removed = Arc::make_mut(&mut self.sessions).remove(key);
276 debug_assert!(removed);
277 removed
278 }
279
280 fn peer_can_send(&self, key: &SessionKey) -> bool {
281 self.sessions.contains(key) && !self.peer_shutdown.contains(key)
282 }
283
284 fn shutdown_peer(&mut self, key: &SessionKey, kind: TransportKind) -> bool {
286 kind == TransportKind::Tcp
287 && self.sessions.contains(key)
288 && !self.peer_shutdown.contains(key)
289 && Arc::make_mut(&mut self.peer_shutdown).insert(key.clone())
290 }
291}
292
293#[derive(Clone)]
295pub struct Relay<T> {
296 namespace: String,
297 kind: TransportKind,
298 config: HashMap<String, T>,
299}
300
301impl<T> Relay<T> {
302 pub fn tcp(config: HashMap<String, T>) -> Self {
304 Self {
305 namespace: TCP.to_string(),
306 kind: TransportKind::Tcp,
307 config,
308 }
309 }
310
311 pub fn udp(config: HashMap<String, T>) -> Self {
313 Self {
314 namespace: UDP.to_string(),
315 kind: TransportKind::Udp,
316 config,
317 }
318 }
319}
320
321impl<T> Protocol for Relay<T>
322where T: Clone + DeserializeOwned + Serialize + MaybeSend + 'static
323{
324 type State = RelayState<T>;
325 type Event = RelayEvent<T>;
326 type Effect = RelayEffect<T>;
327
328 fn namespace(&self) -> &str {
329 self.namespace.as_str()
330 }
331
332 fn init(&self) -> RelayState<T> {
333 RelayState {
334 services: Arc::new(self.config.clone()),
335 sessions: Arc::new(HashSet::new()),
336 session_quota: Arc::new(PeerQuota::new(
337 MAX_RELAY_SESSIONS,
338 MAX_RELAY_SESSIONS_PER_PEER,
339 )),
340 peer_shutdown: Arc::new(HashSet::new()),
341 next_session: 0,
342 }
343 }
344
345 fn decode(&self, wire: Wire<'_>) -> Result<RelayEvent<T>, Reject> {
346 if wire.from == wire.me {
347 let command = rings_codec::deserialize::<RelayCommand<T>>(wire.payload)
348 .map_err(|e| Reject(format!("bad relay command: {e}")))?;
349 Ok(RelayEvent::Command(command))
350 } else {
351 let frame = rings_codec::deserialize::<Frame>(wire.payload)
352 .map_err(|e| Reject(format!("bad relay frame: {e}")))?;
353 Ok(RelayEvent::Frame {
354 from: wire.from,
355 frame,
356 })
357 }
358 }
359
360 fn step(
361 &self,
362 ctx: Ctx<'_, RelayState<T>>,
363 event: RelayEvent<T>,
364 ) -> Transition<RelayState<T>, RelayEffect<T>> {
365 match event {
366 RelayEvent::Command(command) => {
367 step_command(self.namespace.as_str(), ctx.state, command)
368 }
369 RelayEvent::Frame { from, frame } => {
370 step_frame(self.kind, self.namespace.as_str(), ctx.state, from, frame)
371 }
372 }
373 }
374}
375
376fn step_command<T: Clone>(
381 namespace: &str,
382 state: &RelayState<T>,
383 command: RelayCommand<T>,
384) -> Transition<RelayState<T>, RelayEffect<T>> {
385 let mut next = state.clone();
386 match command {
387 RelayCommand::RegisterService { name, target } => {
388 Arc::make_mut(&mut next.services).insert(name, target);
389 Transition::pure(next)
390 }
391 RelayCommand::Accepted {
392 token,
393 peer,
394 service,
395 } => {
396 let session = SessionId(next.next_session);
399 let key = SessionKey::new(peer, namespace, session, Initiator::Local);
400 if !next.can_admit_session(&key) {
401 return Transition::with(next, vec![RelayEffect::RejectAccepted { token }]);
402 }
403 let Some(next_session) = next.next_session.checked_add(1) else {
404 return Transition::with(next, vec![RelayEffect::RejectAccepted { token }]);
405 };
406 next.next_session = next_session;
407 if !next.insert_session(key.clone()) {
409 return Transition::with(next, vec![RelayEffect::RejectAccepted { token }]);
410 }
411 Transition::with(next, vec![RelayEffect::OpenAccepted {
412 token,
413 key,
414 service,
415 }])
416 }
417 RelayCommand::Untrack {
418 peer,
419 session,
420 initiator,
421 } => {
422 next.remove_session(&SessionKey::new(peer, namespace, session, initiator));
423 Transition::pure(next)
424 }
425 RelayCommand::Abort {
426 peer,
427 session,
428 initiator,
429 } => {
430 let key = SessionKey::new(peer, namespace, session, initiator);
431 if next.remove_session(&key) {
432 Transition::with(next, vec![RelayEffect::SendClose {
433 to: peer,
434 session,
435 from_opener: matches!(initiator, Initiator::Local),
436 }])
437 } else {
438 Transition::pure(next)
439 }
440 }
441 }
442}
443
444fn step_frame<T: Clone>(
446 kind: TransportKind,
447 namespace: &str,
448 state: &RelayState<T>,
449 from: Did,
450 frame: Frame,
451) -> Transition<RelayState<T>, RelayEffect<T>> {
452 match frame {
453 Frame::Open { session, service } => {
455 let key = SessionKey::new(from, namespace, session, Initiator::Remote);
456 if state.sessions.contains(&key) {
460 return Transition::pure(state.clone());
461 }
462 match state.services.get(service.as_str()) {
463 Some(target) => {
464 let mut next = state.clone();
465 if !next.can_admit_session(&key) {
466 return rejected_open(next, key);
467 }
468 let target = target.clone();
469 if !next.insert_session(key.clone()) {
470 return rejected_open(next, key);
471 }
472 Transition::with(next, vec![RelayEffect::Connect { key, target, kind }])
473 }
474 None => rejected_open(state.clone(), key),
475 }
476 }
477 Frame::Data {
481 session,
482 from_opener,
483 bytes,
484 } => {
485 let key = SessionKey::new(from, namespace, session, opener_to_initiator(from_opener));
486 if state.peer_can_send(&key) {
487 Transition::with(state.clone(), vec![RelayEffect::Write { key, bytes }])
488 } else {
489 Transition::pure(state.clone())
490 }
491 }
492 Frame::Shutdown {
493 session,
494 from_opener,
495 } => {
496 let key = SessionKey::new(from, namespace, session, opener_to_initiator(from_opener));
497 let mut next = state.clone();
498 if next.shutdown_peer(&key, kind) {
499 Transition::with(next, vec![RelayEffect::Shutdown { key }])
500 } else {
501 Transition::pure(next)
502 }
503 }
504 Frame::Close {
505 session,
506 from_opener,
507 } => {
508 let key = SessionKey::new(from, namespace, session, opener_to_initiator(from_opener));
509 if state.sessions.contains(&key) {
510 let mut next = state.clone();
511 next.remove_session(&key);
512 Transition::with(next, vec![RelayEffect::Close { key }])
513 } else {
514 Transition::pure(state.clone())
515 }
516 }
517 }
518}
519
520fn rejected_open<T>(
522 state: RelayState<T>,
523 key: SessionKey,
524) -> Transition<RelayState<T>, RelayEffect<T>> {
525 Transition::with(state, vec![RelayEffect::SendClose {
526 to: key.peer,
527 session: key.session,
528 from_opener: false,
530 }])
531}
532
533fn opener_to_initiator(from_opener: bool) -> Initiator {
535 if from_opener {
536 Initiator::Remote
537 } else {
538 Initiator::Local
539 }
540}
541
542pub(crate) fn close_frame(session: SessionId, from_opener: bool) -> crate::error::Result<Bytes> {
545 let frame = Frame::Close {
546 session,
547 from_opener,
548 };
549 rings_codec::serialize(&frame)
550 .map(Bytes::from)
551 .map_err(|_| crate::error::Error::EncodeError)
552}
553
554#[cfg(rings_native)]
560pub(crate) struct NativeRelay {
561 engine: Arc<crate::extension::transport::engine::TransportSessions>,
562 control_outbox: ControlOutbox,
563}
564
565#[cfg(rings_native)]
566impl NativeRelay {
567 pub(crate) fn new(engine: Arc<crate::extension::transport::engine::TransportSessions>) -> Self {
569 Self {
570 engine,
571 control_outbox: ControlOutbox::default(),
572 }
573 }
574
575 #[cfg(all(test, rings_native))]
576 pub(crate) fn new_with_control_send_test_hook(
577 engine: Arc<crate::extension::transport::engine::TransportSessions>,
578 hook: Arc<ControlSendTestHook>,
579 ) -> Self {
580 Self {
581 engine,
582 control_outbox: ControlOutbox::with_test_hook(hook),
583 }
584 }
585}
586
587#[cfg(rings_native)]
588#[async_trait::async_trait]
589impl Interpret for NativeRelay {
590 type Effect = RelayEffect<std::net::SocketAddr>;
591
592 async fn run(
593 &self,
594 scope: &EffectScope,
595 effect: RelayEffect<std::net::SocketAddr>,
596 ) -> crate::error::Result<Vec<Bytes>> {
597 match effect {
598 RelayEffect::Connect { key, target, kind } => {
599 let admission =
600 self.engine
601 .clone()
602 .connect(scope.lifecycle(), key.clone(), target, kind);
603 return enqueue_feedback::<std::net::SocketAddr>(key, admission);
604 }
605 RelayEffect::Write { key, bytes } => {
606 let admission = self.engine.write(&key, bytes);
607 return enqueue_feedback::<std::net::SocketAddr>(key, admission);
608 }
609 RelayEffect::Shutdown { key } => {
610 let admission = self.engine.shutdown(&key);
611 return enqueue_feedback::<std::net::SocketAddr>(key, admission);
612 }
613 RelayEffect::Close { key } => {
614 self.engine.close_for_effect(&key);
615 }
616 RelayEffect::SendClose {
617 to,
618 session,
619 from_opener,
620 } => {
621 self.control_outbox.enqueue(
622 scope.lifecycle(),
623 to,
624 close_frame(session, from_opener)?,
625 )?;
626 }
627 RelayEffect::OpenAccepted {
628 token,
629 key,
630 service,
631 } => {
632 let feedback =
633 self.engine
634 .clone()
635 .bind_accepted(scope.lifecycle(), token, key, service);
636 return feedback
637 .map(untrack_feedback::<std::net::SocketAddr>)
638 .transpose()
639 .map(|feedback| feedback.into_iter().collect());
640 }
641 RelayEffect::RejectAccepted { token } => {
642 self.engine.evict_pending_for_effect(token);
643 }
644 }
645 Ok(Vec::new())
646 }
647}
648
649fn untrack_feedback<T: Serialize>(key: SessionKey) -> crate::error::Result<Bytes> {
651 let command = RelayCommand::<T>::Untrack {
652 peer: key.peer,
653 session: key.session,
654 initiator: key.initiator,
655 };
656 rings_codec::serialize(&command)
657 .map(Bytes::from)
658 .map_err(|_| crate::error::Error::EncodeError)
659}
660
661fn abort_feedback<T: Serialize>(key: SessionKey) -> crate::error::Result<Bytes> {
664 let command = RelayCommand::<T>::Abort {
665 peer: key.peer,
666 session: key.session,
667 initiator: key.initiator,
668 };
669 rings_codec::serialize(&command)
670 .map(Bytes::from)
671 .map_err(|_| crate::error::Error::EncodeError)
672}
673
674fn enqueue_feedback<T: Serialize>(
676 key: SessionKey,
677 admission: EffectEnqueue,
678) -> crate::error::Result<Vec<Bytes>> {
679 match admission {
680 EffectEnqueue::Enqueued => Ok(Vec::new()),
681 EffectEnqueue::Missing => untrack_feedback::<T>(key).map(|feedback| vec![feedback]),
682 EffectEnqueue::Failed => abort_feedback::<T>(key).map(|feedback| vec![feedback]),
683 }
684}
685
686#[cfg(rings_browser)]
690pub(crate) struct WtRelay {
691 engine: Arc<crate::extension::transport::wt::WtSessions>,
692 control_outbox: ControlOutbox,
693}
694
695#[cfg(rings_browser)]
696impl WtRelay {
697 pub(crate) fn new(engine: Arc<crate::extension::transport::wt::WtSessions>) -> Self {
699 Self {
700 engine,
701 control_outbox: ControlOutbox::default(),
702 }
703 }
704}
705
706#[cfg(rings_browser)]
707#[async_trait::async_trait(?Send)]
708impl Interpret for WtRelay {
709 type Effect = RelayEffect<String>;
710
711 async fn run(
712 &self,
713 scope: &EffectScope,
714 effect: RelayEffect<String>,
715 ) -> crate::error::Result<Vec<Bytes>> {
716 match effect {
717 RelayEffect::Connect { key, target, kind } => {
718 let admission =
719 self.engine
720 .clone()
721 .connect(scope.lifecycle(), key.clone(), target, kind);
722 return enqueue_feedback::<String>(key, admission);
723 }
724 RelayEffect::Write { key, bytes } => {
725 let admission = self.engine.write(scope.lifecycle(), key.clone(), bytes);
726 return enqueue_feedback::<String>(key, admission);
727 }
728 RelayEffect::Shutdown { key } => {
729 let admission = self.engine.shutdown(scope.lifecycle(), key.clone());
730 return enqueue_feedback::<String>(key, admission);
731 }
732 RelayEffect::Close { key } => {
733 self.engine.close_for_effect(&key);
734 }
735 RelayEffect::SendClose {
736 to,
737 session,
738 from_opener,
739 } => {
740 self.control_outbox.enqueue(
741 scope.lifecycle(),
742 to,
743 close_frame(session, from_opener)?,
744 )?;
745 }
746 RelayEffect::OpenAccepted { .. } => {
749 tracing::warn!("browser relay received OpenAccepted; it has no local listener");
750 }
751 RelayEffect::RejectAccepted { .. } => {
752 tracing::warn!("browser relay received RejectAccepted; it has no local listener");
753 }
754 }
755 Ok(Vec::new())
756 }
757}
758
759#[cfg(rings_native)]
769#[derive(Clone)]
770pub struct RelayHandle {
771 engine: Arc<crate::extension::transport::engine::TransportSessions>,
772 tcp: Scope,
773 udp: Scope,
774}
775
776#[cfg(rings_native)]
777impl RelayHandle {
778 pub fn install(extensions: &crate::extension::ext::Extensions) -> crate::error::Result<Self> {
783 let engine = Arc::new(crate::extension::transport::engine::TransportSessions::new());
784 extensions.register_many(vec![
786 (Relay::tcp(HashMap::new()), NativeRelay::new(engine.clone())),
787 (Relay::udp(HashMap::new()), NativeRelay::new(engine.clone())),
788 ])?;
789 let core = extensions.core();
790 Ok(Self {
791 engine,
792 tcp: Scope::new(core.clone(), TCP.to_string()),
793 udp: Scope::new(core, UDP.to_string()),
794 })
795 }
796
797 pub async fn open_tcp_tunnel(
800 &self,
801 local_addr: std::net::SocketAddr,
802 peer: Did,
803 service: String,
804 ) -> crate::error::Result<()> {
805 self.open_tunnel(&self.tcp, local_addr, peer, service, TransportKind::Tcp)
806 .await
807 }
808
809 pub async fn relay_tcp_stream(
811 &self,
812 stream: tokio::net::TcpStream,
813 peer: Did,
814 service: String,
815 ) -> crate::error::Result<()> {
816 self.engine
817 .clone()
818 .relay_tcp_stream(self.tcp.clone(), stream, peer, service)
819 .await;
820 Ok(())
821 }
822
823 pub async fn open_udp_tunnel(
826 &self,
827 local_addr: std::net::SocketAddr,
828 peer: Did,
829 service: String,
830 ) -> crate::error::Result<()> {
831 self.open_tunnel(&self.udp, local_addr, peer, service, TransportKind::Udp)
832 .await
833 }
834
835 async fn open_tunnel(
836 &self,
837 scope: &Scope,
838 local_addr: std::net::SocketAddr,
839 peer: Did,
840 service: String,
841 kind: TransportKind,
842 ) -> crate::error::Result<()> {
843 self.engine
847 .clone()
848 .listen(scope.clone(), local_addr, peer, service, kind)
849 .await;
850 Ok(())
851 }
852
853 pub async fn register_tcp_service(
855 &self,
856 name: String,
857 addr: std::net::SocketAddr,
858 ) -> crate::error::Result<()> {
859 register_service(&self.tcp, name, addr).await
860 }
861
862 pub async fn register_udp_service(
864 &self,
865 name: String,
866 addr: std::net::SocketAddr,
867 ) -> crate::error::Result<()> {
868 register_service(&self.udp, name, addr).await
869 }
870}
871
872#[cfg(any(rings_native, rings_browser))]
875async fn register_service<T>(scope: &Scope, name: String, target: T) -> crate::error::Result<()>
876where T: Serialize {
877 let command = RelayCommand::RegisterService { name, target };
878 let payload = rings_codec::serialize(&command).map_err(|_| crate::error::Error::EncodeError)?;
879 scope.inject(Bytes::from(payload)).await
880}
881
882#[cfg(rings_browser)]
887#[derive(Clone)]
888pub struct RelayHandle {
889 tcp: Scope,
890 udp: Scope,
891}
892
893#[cfg(rings_browser)]
894impl RelayHandle {
895 pub fn install(extensions: &crate::extension::ext::Extensions) -> crate::error::Result<Self> {
905 let engine = Arc::new(crate::extension::transport::wt::WtSessions::new());
906 extensions.register_many(vec![
908 (Relay::tcp(HashMap::new()), WtRelay::new(engine.clone())),
909 (Relay::udp(HashMap::new()), WtRelay::new(engine)),
910 ])?;
911 let core = extensions.core();
912 Ok(Self {
913 tcp: Scope::new(core.clone(), TCP.to_string()),
914 udp: Scope::new(core, UDP.to_string()),
915 })
916 }
917
918 pub async fn register_wt_service(&self, name: String, url: String) -> crate::error::Result<()> {
921 register_service(&self.tcp, name, url).await
922 }
923
924 pub async fn register_wt_udp_service(
927 &self,
928 name: String,
929 url: String,
930 ) -> crate::error::Result<()> {
931 register_service(&self.udp, name, url).await
932 }
933}
934
935#[cfg(test)]
936mod test_relay;