Skip to main content

rings_node/extension/protocols/
relay.rs

1//! Generic transport-relay protocol — one pure server-side state machine for TCP and UDP,
2//! native and browser.
3//!
4//! The pure model is generic over the **target** `T` a service resolves to: a
5//! `SocketAddr` natively, a WebTransport `Url` (string) in the browser. The same `step`,
6//! state, duplicate-`Open` rejection and owner-rejection serve both — only the
7//! *interpreter* differs (native `NativeRelay` over OS sockets, browser `WtRelay` over
8//! WebTransport). This is the code realization of "TCP/UDP/native/browser are one relay".
9//!
10//! Every session is identified by the **owner-scoped key** `(from, namespace, session,
11//! initiator)` ([`SessionKey`]). `from` is the authenticated sender (owner rejection: a peer
12//! can only name keys whose `from` is itself), and `initiator` records which end opened it —
13//! so a session a peer opened never collides with one we opened that got the same id
14//! (bidirectional-open safety). A frame's `from_opener` flips to our `initiator`.
15//!
16//! The reducer is the **sole authority** over the session set: `Data`/`Shutdown`/`Close`
17//! emit an effect only for a session in `sessions` (the engine never adjudicates liveness).
18//!
19//! ```text
20//!   S = (services : Name ⇀ T, sessions : ℘ SessionKey, next : ℕ)
21//!   k = (from, namespace, session, init)        init = Remote if from_opener else Local
22//!   step (Command(Register n t))                ↦ (S[services∪{n↦t}], ε)
23//!   step (Command(Accepted tok peer svc))       ↦ (S[sessions∪{kₗ}, next+1], [OpenAccepted tok kₗ svc])
24//!                                                   where kₗ=(peer,ns,next,Local)   ← core mints the id
25//!   step (Command(Untrack k))                   ↦ (S[sessions∖{k}], ε)
26//!   step (Frame(from, Open s n)) | k∈sessions   ↦ (S, ε)                  (live duplicate)
27//!                                | n∈services    ↦ (S∪{k}, [Connect k t kind])
28//!                                | otherwise     ↦ (S, [SendClose s])
29//!   step (Frame(from, Data s b)) | k∈sessions∖shutdown ↦ (S, [Write k b]) else (S, ε)
30//!   step (Frame(from, FIN s))    | TCP ∧ k∈sessions    ↦ (S∪shutdown(k), [Shutdown k])
31//!                                                   repeated/UDP ↦ (S, ε)
32//!   step (Frame(from, Close s))  | k∈sessions          ↦ (S∖{k}, [Close k]) else (S, ε)
33//!   invariant                    |sessions| ≤ 1024 ∧ ∀peer. sessions(peer) ≤ 64
34//! ```
35
36use 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
72/// Namespace for the TCP relay.
73pub const TCP: &str = "tcp";
74/// Namespace for the UDP relay.
75pub const UDP: &str = "udp";
76
77/// Hard per-namespace live-session bound. The reducer owns admission, so both native and
78/// browser interpreters inherit the same resource ceiling.
79pub(crate) const MAX_RELAY_SESSIONS: usize = 1_024;
80/// A single authenticated peer cannot consume the entire relay-session budget.
81pub(crate) const MAX_RELAY_SESSIONS_PER_PEER: usize = 64;
82/// A local control command, re-injected by the provider (provenance = self; never sent by
83/// peers). Generic over the service target `T`.
84#[derive(Clone, Debug, Serialize, Deserialize, PartialEq, Eq)]
85pub enum RelayCommand<T> {
86    /// Map a service name to a local target that `Open` may dial.
87    RegisterService {
88        /// Service name.
89        name: String,
90        /// Local target (`SocketAddr` natively, WebTransport URL in browser).
91        target: T,
92    },
93    /// Engine→protocol feedback: a local connection/datagram-flow was accepted, pending
94    /// under engine-local `token`, destined for `peer`'s `service`. The pure `step` mints
95    /// the session id (so id allocation lives in the core, not the shell), records it, and
96    /// replies with [`RelayEffect::OpenAccepted`] to bind the pending resource. The engine
97    /// never mints or decides identity — it only reports the raw accept and executes effects.
98    Accepted {
99        /// Engine-local handle for the pending (not-yet-bound) connection/flow.
100        token: u64,
101        /// The remote peer this session is tunnelled to.
102        peer: Did,
103        /// The remote service to open.
104        service: String,
105    },
106    /// Engine→protocol feedback: a session was torn down by the engine (any side); forget
107    /// it. The single point through which every teardown reaches the pure state.
108    Untrack {
109        /// The remote peer of the session.
110        peer: Did,
111        /// The session id.
112        session: SessionId,
113        /// Which end opened it (so the right key is removed).
114        initiator: Initiator,
115    },
116    /// Engine→protocol feedback: the current backend failed under an ordered
117    /// effect. Forget it and emit exactly one peer-facing `Close`.
118    Abort {
119        /// The remote peer of the session.
120        peer: Did,
121        /// The session id.
122        session: SessionId,
123        /// Which end opened it.
124        initiator: Initiator,
125    },
126}
127
128/// The relay's typed input: a self-injected [`RelayCommand`] or an authenticated peer
129/// [`Frame`]. The `from == me` split is resolved in [`Relay::decode`].
130pub enum RelayEvent<T> {
131    /// Runtime service registration (provenance = self).
132    Command(RelayCommand<T>),
133    /// A network frame from an authenticated peer.
134    Frame {
135        /// Authenticated sender.
136        from: Did,
137        /// The frame.
138        frame: Frame,
139    },
140}
141
142/// The relay's own effect algebra (interpreted by `NativeRelay` / `WtRelay`).
143#[derive(Clone, Debug, PartialEq, Eq)]
144pub enum RelayEffect<T> {
145    /// Open a local backend session to `target` and relay it (peer opened a session).
146    Connect {
147        /// Owner-scoped session key.
148        key: SessionKey,
149        /// Local target to dial.
150        target: T,
151        /// Stream (TCP) or datagram (UDP).
152        kind: TransportKind,
153    },
154    /// Write peer bytes to a session's local stream.
155    Write {
156        /// Target session.
157        key: SessionKey,
158        /// Bytes.
159        bytes: Bytes,
160    },
161    /// Half-close a session's local write side (peer FIN).
162    Shutdown {
163        /// Target session.
164        key: SessionKey,
165    },
166    /// Close a session (full teardown).
167    Close {
168        /// Target session.
169        key: SessionKey,
170    },
171    /// Reply a `Frame::Close` to a peer that opened an unknown service. The reply goes out under
172    /// the interpreter's own namespace (its [`Scope`]), so the effect carries no namespace of
173    /// its own.
174    SendClose {
175        /// Peer to reply to.
176        to: Did,
177        /// Session id to close.
178        session: SessionId,
179        /// Whether *we* opened the session (false: the peer did).
180        from_opener: bool,
181    },
182    /// Bind a pending accepted connection/flow (engine-local `token`) to the session `key`
183    /// the pure `step` just minted, then open it to the peer and start relaying. The reply
184    /// to [`RelayCommand::Accepted`] — this is how a step-minted id reaches the engine.
185    OpenAccepted {
186        /// Engine-local handle for the pending connection/flow.
187        token: u64,
188        /// The session key minted by the pure step.
189        key: SessionKey,
190        /// The remote service to open.
191        service: String,
192    },
193    /// Drop an accepted local resource when the pure session-id space is exhausted.
194    RejectAccepted {
195        /// Engine-local handle whose resource must be released.
196        token: u64,
197    },
198}
199
200/// Relay state: the service registry and the set of live sessions in both directions. The live
201/// OS/WebTransport resources are the interpreter's engine table; this is the protocol's view
202/// used for admission, ownership checks, and duplicate-`Open` rejection.
203#[derive(Clone)]
204pub struct RelayState<T> {
205    services: Arc<HashMap<String, T>>,
206    /// Persistent live-session index. Cloning a pure state for a data-frame transition is O(1);
207    /// session lifecycle steps copy on write.
208    sessions: Arc<HashSet<SessionKey>>,
209    /// Exact cached cardinality of `sessions` projected by authenticated peer.
210    ///
211    /// Invariant: `session_quota.peer_total(p) = |{ k \in sessions : k.peer = p }|`; zero entries are
212    /// absent. Keeping the projection in the pure state makes admission O(1) without giving the
213    /// interpreter a second source of truth.
214    session_quota: Arc<PeerQuota>,
215    /// TCP sessions whose peer-to-local direction consumed its affine FIN.
216    ///
217    /// Invariant: `peer_shutdown` is a subset of `sessions`. The reverse direction remains live
218    /// until `Close`, but no later peer `Data` can escape the pure reducer into a backend writer.
219    peer_shutdown: Arc<HashSet<SessionKey>>,
220    /// Monotonic allocator for client-side session ids. Lives in the **pure** state so the
221    /// core (not the engine) mints session identities — `Event → step → Effect` is the sole
222    /// authority for both the session set and its ids.
223    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    /// Consume the peer-to-local FIN exactly once for a live TCP stream.
285    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/// Transport relay protocol (server side), generic over the service target `T`.
294#[derive(Clone)]
295pub struct Relay<T> {
296    namespace: String,
297    kind: TransportKind,
298    config: HashMap<String, T>,
299}
300
301impl<T> Relay<T> {
302    /// A TCP relay with a fixed service configuration.
303    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    /// A UDP relay with a fixed service configuration.
312    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
376/// Apply a local [`RelayCommand`]. Pure. `Accepted`/`Untrack` are the engine→protocol
377/// feedback that make `step` the sole authority over the session set **and its ids**: the
378/// core mints the id on `Accepted` (the engine reported only a local token) and forgets the
379/// session on `Untrack`.
380fn 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            // The core mints the session id (the engine reported only its local token), so
397            // id allocation is part of the pure state transition, not a shell decision.
398            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            // A locally-accepted tunnel: we are the initiator.
408            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
444/// Apply a network [`Frame`]. Pure; emits relay effects scoped to the authenticated `from`.
445fn 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        // `Open` is always sent by the opener, so from our side the peer is the initiator.
454        Frame::Open { session, service } => {
455            let key = SessionKey::new(from, namespace, session, Initiator::Remote);
456            // A duplicate for a live session is silent. Rejected opens are not retained in pure
457            // state: the bounded interpreter outbox owns control-plane resource admission, so a
458            // historical failure cannot permanently consume protocol capacity.
459            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        // Data/Shutdown/Close are guarded on the authoritative session set: the *reducer*
478        // decides whether the effect happens, not the engine table. `from_opener` (the
479        // sender opened it) flips to our initiator.
480        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
520/// Reject a peer-opened session without allocating live-session or interpreter state.
521fn 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        // The peer opened it; this endpoint did not.
529        from_opener: false,
530    }])
531}
532
533/// Map a frame's `from_opener` (the **sender** opened the session) to our own [`Initiator`].
534fn opener_to_initiator(from_opener: bool) -> Initiator {
535    if from_opener {
536        Initiator::Remote
537    } else {
538        Initiator::Local
539    }
540}
541
542/// Encode a `Frame::Close` as bytes for an overlay send. `from_opener` is whether *we* (the
543/// sender of this close) opened the session.
544pub(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// ── Native interpreter (OS sockets) ───────────────────────────────────────────────────
555
556/// Native relay interpreter: runs [`RelayEffect`]s over the OS-socket engine it owns. The
557/// engine uses the namespace-scoped [`Scope`] capability for both overlay sends and lifecycle
558/// feedback (`Accepted`/`Untrack`), so the engine has no `Processor` of its own.
559#[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    /// Build over a shared engine.
568    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
649/// Encode an engine teardown as the relay's synchronous, ordered feedback path.
650fn 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
661/// Encode a synchronous engine failure so the pure reducer owns both state
662/// removal and the exactly-once peer-facing terminal effect.
663fn 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
674/// Map backend admission into the relay reducer's ordered feedback algebra.
675fn 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// ── Browser interpreter (WebTransport) ────────────────────────────────────────────────
687
688/// Browser relay interpreter: runs [`RelayEffect`]s over the WebTransport engine it owns.
689#[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    /// Build over a shared WebTransport engine.
698    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            // The browser relay is server-side only (no local listener), so it never reports
747            // an `Accepted` and thus never receives `OpenAccepted`.
748            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// ── Client-side relay handle ──────────────────────────────────────────────────────────
760
761/// Client-facing handle to the relay extension's live engine: open local tunnels and register
762/// local services. This is the relay extension's *own* surface — the relay owns its engine and
763/// installs itself ([`install`](RelayHandle::install)), so nothing about it leaks into the
764/// generic [`Provider`](crate::provider::Provider).
765/// Cloneable; every clone drives the same shared engine and pure [`Relay`] state.
766/// Holds the two per-namespace scoped capabilities (`tcp` / `udp`); each method picks one and
767/// can only act within it, so the handle cannot address an arbitrary namespace even internally.
768#[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    /// Install the relay into an extension registry: register the TCP and UDP interpreters
779    /// over a fresh, relay-owned OS-socket engine and return the client handle. Errors if the
780    /// `tcp`/`udp` namespaces are already taken. Call once per node, after constructing the
781    /// provider — the relay is opt-in, not a `Provider` invariant.
782    pub fn install(extensions: &crate::extension::ext::Extensions) -> crate::error::Result<Self> {
783        let engine = Arc::new(crate::extension::transport::engine::TransportSessions::new());
784        // Atomic: both namespaces register together, or neither (no half-installed relay).
785        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    /// Open a local **TCP** tunnel: bind `local_addr` and relay each accepted connection to
798    /// `peer`'s `service` (client side, forward proxy).
799    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    /// Relay one already-accepted **TCP** stream to `peer`'s `service`.
810    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    /// Open a local **UDP** tunnel: bind `local_addr` and relay each datagram flow to `peer`'s
824    /// `service` (client side, forward proxy).
825    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        // Bind a local listener on the relay engine with this namespace's scope. Each accepted
844        // connection is reported back through the pure relay (`Accepted`), so
845        // `RelayState.sessions` stays the sole authority.
846        self.engine
847            .clone()
848            .listen(scope.clone(), local_addr, peer, service, kind)
849            .await;
850        Ok(())
851    }
852
853    /// Register (at runtime) a local service the **TCP** relay may dial (`name` → `addr`).
854    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    /// Register (at runtime) a local service the **UDP** relay may dial (`name` → `addr`).
863    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/// Map a service `name` → `target` by self-injecting a `RegisterService` command into the
873/// scope's own namespace (provenance = self).
874#[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/// Client-facing handle to the browser relay extension's live WebTransport engine. It owns the
883/// two per-namespace scoped capabilities (`tcp` / `udp`) and registers services, but exposes no
884/// tunnel-open surface because the browser relay is server-side only. Cloneable; see the native
885/// [`RelayHandle`].
886#[cfg(rings_browser)]
887#[derive(Clone)]
888pub struct RelayHandle {
889    tcp: Scope,
890    udp: Scope,
891}
892
893#[cfg(rings_browser)]
894impl RelayHandle {
895    /// Install the browser relay into an extension registry: register the TCP and UDP
896    /// interpreters over a fresh, relay-owned WebTransport engine and return the client handle.
897    /// Errors if the `tcp`/`udp` namespaces are already taken. Call once per node, after
898    /// constructing the provider — the relay is opt-in, not a `Provider` invariant.
899    ///
900    /// This is a **Rust-wasm-facing** surface: there is no `wasm_bindgen` install/handle for JS
901    /// yet (unlike `provider.on(...)`), so browser relay is reachable only from Rust-wasm apps.
902    /// A JS-facing extension install API can be added when a JS consumer needs WebTransport
903    /// relay; it must not put these methods back on the generic `Provider`.
904    pub fn install(extensions: &crate::extension::ext::Extensions) -> crate::error::Result<Self> {
905        let engine = Arc::new(crate::extension::transport::wt::WtSessions::new());
906        // Atomic: both namespaces register together, or neither (no half-installed relay).
907        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    /// Register a WebTransport-backed service for the browser **TCP** relay, mapping
919    /// `name` → WebTransport `url` (under the `tcp` namespace).
920    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    /// Register a WebTransport-backed service for the browser **UDP** relay (datagrams),
925    /// mapping `name` → WebTransport `url` (under the `udp` namespace).
926    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;