Skip to main content

rings_node/onion/tcp/
mod.rs

1//! Native TCP adapter for route-aware onion circuits.
2
3use std::collections::hash_map::Entry;
4use std::collections::HashMap;
5use std::sync::Arc;
6use std::sync::Mutex;
7use std::time::Duration;
8
9use bytes::Bytes;
10use rings_core::dht::Did;
11use rings_core::ecc::PublicKey;
12use rings_core::session::SessionSk;
13use serde::Deserialize;
14use serde::Serialize;
15use tokio::net::TcpStream;
16use tokio::sync::mpsc;
17use tokio::sync::oneshot;
18use tokio::time::timeout;
19use tokio::time::Instant;
20
21use crate::error::Error;
22use crate::error::Result;
23use crate::extension::ext::Extensions;
24use crate::extension::ext::Scope;
25use crate::onion::circuit::route_first_hop;
26use crate::onion::circuit::send_backward;
27use crate::onion::circuit::OnionAuthenticatedPayload;
28use crate::onion::circuit::OnionBackwardSequence;
29use crate::onion::circuit::OnionCircuitCapabilities;
30use crate::onion::circuit::OnionCircuitExitFrame;
31use crate::onion::circuit::OnionCircuitHandler;
32use crate::onion::circuit::OnionCircuitId;
33use crate::onion::circuit::OnionCircuitPath;
34use crate::onion::circuit::OnionCircuitPayload;
35use crate::onion::circuit::OnionCircuitProtocol;
36use crate::onion::circuit::OnionCircuitShell;
37use crate::onion::circuit::OnionClientReturn;
38use crate::onion::circuit::OnionForwardNonce;
39use crate::onion::circuit::OnionForwardSequence;
40use crate::onion::circuit::OnionLinkSender;
41use crate::onion::circuit::OnionReturnId;
42use crate::onion::circuit::ONION_CIRCUIT_NAMESPACE;
43use crate::onion::exit_accounting::OnionExitAccounting;
44use crate::onion::exit_accounting::OnionExitLease;
45use crate::onion::https::try_handle_https_exit_payload;
46use crate::onion::https::OnionHttpsRuntime;
47use crate::onion::replay::OnionForwardReplayKey;
48use crate::onion::replay::OnionForwardReplayPartitions;
49use crate::onion::replay::OnionSequenceWindow;
50use crate::onion::replay::ReplayAdmission;
51use crate::onion::replay::SequenceAdmission;
52use crate::onion::OnionExitDescriptor;
53use crate::onion::OnionExitFailure;
54use crate::onion::OnionExitPolicy;
55use crate::onion::OnionProxyTarget;
56use crate::onion::OnionRoute;
57use crate::onion::OnionRouteError;
58use crate::onion::OnionServiceName;
59use crate::sync_lock::lock;
60
61mod client;
62mod config;
63mod duplex;
64mod exit;
65mod inbound;
66mod pump;
67
68use client::spawn_client_stream;
69use client::TcpBackwardRoute;
70pub use config::NativeOnionTcpExitConfig;
71#[cfg(test)]
72use duplex::TcpDuplexState;
73use exit::admit_exit_target;
74use exit::connect_exit_target;
75use exit::open_response_deadline;
76use exit::spawn_exit_stream;
77use exit::ExitStreamTask;
78use inbound::TcpInbound;
79
80const TCP_BUF: usize = 30_000;
81const TCP_OPEN_TIMEOUT_SECS: u64 = 30;
82
83#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
84enum OnionTcpPayload {
85    Open { target: String },
86    Opened,
87    Data { bytes: Bytes },
88    Shutdown,
89    Close,
90    Error(OnionExitFailure),
91}
92
93fn encode_tcp_payload(
94    service: &OnionServiceName,
95    payload: OnionTcpPayload,
96) -> Result<OnionCircuitPayload> {
97    rings_codec::serialize(&payload)
98        .map(|body| OnionCircuitPayload::new(service.clone(), Bytes::from(body)))
99        .map_err(|_| Error::EncodeError)
100}
101
102fn decode_tcp_payload_for_service(
103    payload: OnionCircuitPayload,
104    service: &OnionServiceName,
105) -> Result<Option<OnionTcpPayload>> {
106    if !payload.is_service(service) {
107        return Ok(None);
108    }
109    rings_codec::deserialize(payload.body.as_ref())
110        .map(Some)
111        .map_err(|_| Error::DecodeError)
112}
113
114/// Native handle for opening TCP streams over route-aware onion circuits.
115#[derive(Clone)]
116pub struct NativeOnionCircuitHandle {
117    runtime: Arc<OnionTcpRuntime>,
118    scope: Scope,
119}
120
121impl NativeOnionCircuitHandle {
122    /// Install the route-aware onion circuit protocol.
123    pub fn install(
124        extensions: &Extensions,
125        session_sk: SessionSk,
126        allow_relay: bool,
127        exit_config: Option<NativeOnionTcpExitConfig>,
128    ) -> Result<Self> {
129        let allow_exit = exit_config.is_some();
130        let (runtime, https) = native_onion_runtimes(session_sk.clone(), exit_config);
131        if let Some(config) = runtime.exit_config.as_ref() {
132            if config.allows_service(&OnionServiceName::https()) {
133                https.set_exit_policy(Some(config.policy().clone()));
134                https.set_native_proxy(config.https_proxy().map(ToString::to_string));
135            }
136        }
137        let capabilities = OnionCircuitCapabilities::from_registration(allow_relay, allow_exit);
138        let handler_session_sk = session_sk.clone();
139        extensions.register(
140            OnionCircuitProtocol::new(capabilities),
141            OnionCircuitShell::with_link_sender(
142                session_sk,
143                NativeOnionCircuitHandler {
144                    runtime: runtime.clone(),
145                    https,
146                    session_sk: handler_session_sk,
147                },
148                runtime.link_sender.clone(),
149            ),
150        )?;
151        Ok(Self {
152            runtime,
153            scope: Scope::new(extensions.core(), ONION_CIRCUIT_NAMESPACE.to_string()),
154        })
155    }
156
157    /// Relay an already-accepted TCP stream over `route`.
158    pub async fn relay_tcp_stream(
159        &self,
160        stream: TcpStream,
161        route: OnionRoute,
162        target: OnionProxyTarget,
163    ) -> Result<()> {
164        let opened = self.open_tcp_stream(route, target).await?;
165        opened.relay(stream);
166        Ok(())
167    }
168
169    /// Open a TCP stream over `route` and wait until the exit has connected the target.
170    pub async fn open_tcp_stream(
171        &self,
172        route: OnionRoute,
173        target: OnionProxyTarget,
174    ) -> Result<NativeOnionOpenStream> {
175        self.runtime
176            .open_client_connection(self.scope.clone(), route, target)
177            .await
178    }
179}
180
181fn native_onion_runtimes(
182    session_sk: SessionSk,
183    exit_config: Option<NativeOnionTcpExitConfig>,
184) -> (Arc<OnionTcpRuntime>, Arc<OnionHttpsRuntime>) {
185    let accounting = OnionExitAccounting::default();
186    let link_sender = OnionLinkSender::default();
187    let runtime = Arc::new(OnionTcpRuntime::with_resources(
188        session_sk,
189        exit_config,
190        accounting.clone(),
191        link_sender.clone(),
192    ));
193    let https = Arc::new(OnionHttpsRuntime::with_resources(accounting, link_sender));
194    (runtime, https)
195}
196
197/// Client-side onion TCP stream after the exit has accepted and connected the target.
198pub struct NativeOnionOpenStream {
199    runtime: Arc<OnionTcpRuntime>,
200    scope: Scope,
201    key: TcpStreamKey,
202    path: OnionCircuitPath,
203    client_return: OnionClientReturn,
204    rx: mpsc::Receiver<TcpInbound>,
205}
206
207impl NativeOnionOpenStream {
208    /// Relay `stream` through this already-open onion TCP stream.
209    pub fn relay<S>(self, stream: S)
210    where S: tokio::io::AsyncRead + tokio::io::AsyncWrite + Unpin + Send + 'static {
211        spawn_client_stream(
212            self.runtime,
213            self.scope,
214            self.key,
215            stream,
216            self.path,
217            self.client_return,
218            self.rx,
219        );
220    }
221}
222
223#[derive(Clone)]
224struct NativeOnionCircuitHandler {
225    runtime: Arc<OnionTcpRuntime>,
226    https: Arc<OnionHttpsRuntime>,
227    session_sk: SessionSk,
228}
229
230#[async_trait::async_trait]
231impl OnionCircuitHandler for NativeOnionCircuitHandler {
232    async fn handle_exit(&self, scope: &Scope, frame: OnionCircuitExitFrame) -> Result<()> {
233        if frame
234            .payload
235            .matches_service(crate::onion::proxy::ONION_PROXY_HTTPS_SERVICE)
236            && try_handle_https_exit_payload(&self.https, &self.session_sk, scope, frame.clone())
237                .await?
238        {
239            return Ok(());
240        }
241        self.runtime.handle_exit_payload(scope.clone(), frame).await
242    }
243
244    async fn handle_client(
245        &self,
246        _scope: &Scope,
247        from: Did,
248        circuit_id: OnionCircuitId,
249        payload: OnionAuthenticatedPayload,
250    ) -> Result<()> {
251        self.runtime
252            .handle_client_payload(from, circuit_id, payload)
253            .await
254    }
255}
256
257struct OnionTcpRuntime {
258    session_sk: SessionSk,
259    client_streams: Mutex<HashMap<TcpStreamKey, ClientStream>>,
260    exit_streams: Mutex<HashMap<TcpStreamKey, ExitStream>>,
261    forward_replays: Mutex<OnionForwardReplayPartitions>,
262    exit_config: Option<NativeOnionTcpExitConfig>,
263    accounting: OnionExitAccounting,
264    link_sender: OnionLinkSender,
265}
266
267impl OnionTcpRuntime {
268    #[cfg(test)]
269    fn new(session_sk: SessionSk, exit_config: Option<NativeOnionTcpExitConfig>) -> Self {
270        Self::with_resources(
271            session_sk,
272            exit_config,
273            OnionExitAccounting::default(),
274            OnionLinkSender::default(),
275        )
276    }
277
278    fn with_resources(
279        session_sk: SessionSk,
280        exit_config: Option<NativeOnionTcpExitConfig>,
281        accounting: OnionExitAccounting,
282        link_sender: OnionLinkSender,
283    ) -> Self {
284        Self {
285            session_sk,
286            client_streams: Mutex::new(HashMap::new()),
287            exit_streams: Mutex::new(HashMap::new()),
288            forward_replays: Mutex::new(OnionForwardReplayPartitions::default()),
289            exit_config,
290            accounting,
291            link_sender,
292        }
293    }
294
295    async fn open_client_connection(
296        self: &Arc<Self>,
297        scope: Scope,
298        route: OnionRoute,
299        target: OnionProxyTarget,
300    ) -> Result<NativeOnionOpenStream> {
301        let expected_return_peer = route_first_hop(&route)?;
302        let expected_exit = route.exit().clone();
303        let service = route.service_name().clone();
304        let client_return = OnionClientReturn::new(self.session_sk.session_public_key());
305        let (tx, rx) = mpsc::channel(32);
306        let (open_tx, open_rx) = oneshot::channel();
307        let key = self.insert_client_stream(
308            service.clone(),
309            expected_return_peer,
310            expected_exit,
311            client_return.return_id,
312            open_tx,
313            tx,
314        )?;
315        let path = match OnionCircuitPath::new(route, key.circuit_id) {
316            Ok(path) => path,
317            Err(error) => {
318                self.remove_client_stream(key);
319                return Err(error);
320            }
321        };
322        let open_payload = match encode_tcp_payload(&service, OnionTcpPayload::Open {
323            target: target.authority(),
324        }) {
325            Ok(payload) => payload,
326            Err(error) => {
327                self.remove_client_stream(key);
328                return Err(error);
329            }
330        };
331        let (first_link, payload) = match path.encode_forward(client_return, open_payload) {
332            Ok(encoded) => encoded,
333            Err(error) => {
334                self.remove_client_stream(key);
335                return Err(error);
336            }
337        };
338        if let Err(error) = self
339            .link_sender
340            .send_sealed(scope.clone(), first_link, payload)
341            .await
342        {
343            self.remove_client_stream(key);
344            return Err(error);
345        }
346        match timeout(Duration::from_secs(TCP_OPEN_TIMEOUT_SECS), open_rx).await {
347            Ok(Ok(Ok(()))) => Ok(NativeOnionOpenStream {
348                runtime: self.clone(),
349                scope,
350                key,
351                path,
352                client_return,
353                rx,
354            }),
355            Ok(Ok(Err(failure))) => {
356                self.remove_client_stream(key);
357                Err(Error::OnionRouteError(OnionRouteError::ExitFailure(
358                    failure,
359                )))
360            }
361            Ok(Err(_)) => {
362                self.remove_client_stream(key);
363                Err(Error::OnionRouteError(
364                    OnionRouteError::TcpOpenResponseClosed,
365                ))
366            }
367            Err(_) => {
368                self.remove_client_stream(key);
369                Err(Error::OnionRouteError(OnionRouteError::TcpOpenTimedOut))
370            }
371        }
372    }
373
374    async fn handle_exit_payload(
375        self: &Arc<Self>,
376        scope: Scope,
377        frame: OnionCircuitExitFrame,
378    ) -> Result<()> {
379        let key = TcpStreamKey {
380            circuit_id: frame.circuit_id,
381        };
382        let Some((service, payload)) = self.decode_exit_payload(frame.payload)? else {
383            return Ok(());
384        };
385        match payload {
386            OnionTcpPayload::Open { target } => {
387                if frame.forward_sequence != OnionForwardSequence::FIRST {
388                    return Err(Error::OnionRouteError(OnionRouteError::ForwardReplay));
389                }
390                self.consume_forward_nonce(frame.from, frame.circuit_id, frame.forward_nonce)?;
391                self.open_exit_stream(TcpExitOpen {
392                    scope,
393                    opened_at: Instant::now(),
394                    key,
395                    circuit_id: frame.circuit_id,
396                    return_peer: frame.return_peer,
397                    return_session_public_key: frame.return_session_public_key,
398                    client: frame.client,
399                    expected_forward_peer: frame.from,
400                    service,
401                    target,
402                })
403                .await
404            }
405            OnionTcpPayload::Data { bytes } => self.send_exit_inbound(
406                key,
407                frame.from,
408                &service,
409                frame.forward_sequence,
410                TcpInbound::Data(bytes),
411            ),
412            OnionTcpPayload::Shutdown => self.send_exit_inbound(
413                key,
414                frame.from,
415                &service,
416                frame.forward_sequence,
417                TcpInbound::Shutdown,
418            ),
419            OnionTcpPayload::Close => self.send_exit_inbound(
420                key,
421                frame.from,
422                &service,
423                frame.forward_sequence,
424                TcpInbound::Close,
425            ),
426            OnionTcpPayload::Opened | OnionTcpPayload::Error(_) => Ok(()),
427        }
428    }
429
430    async fn handle_client_payload(
431        self: &Arc<Self>,
432        from: Did,
433        circuit_id: OnionCircuitId,
434        payload: OnionAuthenticatedPayload,
435    ) -> Result<()> {
436        let key = TcpStreamKey { circuit_id };
437        let payload = self.verify_client_payload(key, from, payload)?;
438        let service = self.client_stream_service(key, from)?;
439        let Some(payload) = decode_tcp_payload_for_service(payload, &service)? else {
440            return Ok(());
441        };
442        match payload {
443            OnionTcpPayload::Data { bytes } => {
444                self.send_client_inbound(key, from, TcpInbound::Data(bytes))
445            }
446            OnionTcpPayload::Shutdown => self.send_client_inbound(key, from, TcpInbound::Shutdown),
447            OnionTcpPayload::Close => self.send_client_inbound(key, from, TcpInbound::Close),
448            OnionTcpPayload::Error(failure) => {
449                if self.complete_client_open(key, from, Err(failure.clone()))? {
450                    return Ok(());
451                }
452                self.send_client_inbound(key, from, TcpInbound::Error(failure))
453            }
454            OnionTcpPayload::Opened => {
455                self.complete_client_open(key, from, Ok(()))?;
456                Ok(())
457            }
458            OnionTcpPayload::Open { .. } => Ok(()),
459        }
460    }
461
462    fn consume_forward_nonce(
463        &self,
464        from: Did,
465        circuit_id: OnionCircuitId,
466        nonce: OnionForwardNonce,
467    ) -> Result<()> {
468        let mut replays = lock(&self.forward_replays)?;
469        match replays.consume(
470            from,
471            OnionForwardReplayKey::new(circuit_id, nonce),
472            rings_core::utils::get_epoch_ms(),
473        ) {
474            ReplayAdmission::Consumed => Ok(()),
475            ReplayAdmission::Duplicate => {
476                Err(Error::OnionRouteError(OnionRouteError::ForwardReplay))
477            }
478            ReplayAdmission::Full => Err(Error::NoPermission),
479        }
480    }
481
482    fn decode_exit_payload(
483        &self,
484        payload: OnionCircuitPayload,
485    ) -> Result<Option<(OnionServiceName, OnionTcpPayload)>> {
486        let service = payload.service_name().clone();
487        if !self.accepts_exit_service(&service) {
488            return Ok(None);
489        }
490        decode_tcp_payload_for_service(payload, &service)
491            .map(|payload| payload.map(|payload| (service, payload)))
492    }
493
494    fn accepts_exit_service(&self, service: &OnionServiceName) -> bool {
495        self.exit_config
496            .as_ref()
497            .is_some_and(|config| config.allows_service(service))
498    }
499
500    async fn open_exit_stream(self: &Arc<Self>, request: TcpExitOpen) -> Result<()> {
501        let Some(exit_config) = &self.exit_config else {
502            return self
503                .reject_exit_open(&request, OnionExitFailure::ExitUnavailable)
504                .await;
505        };
506        if !exit_config.allows_service(&request.service) {
507            return self
508                .reject_exit_open(&request, OnionExitFailure::ExitUnavailable)
509                .await;
510        }
511        let policy = exit_config.policy();
512
513        let target = match admit_exit_target(policy, &request.target) {
514            Ok(target) => target,
515            Err(failure) => return self.reject_exit_open(&request, failure).await,
516        };
517        let (rx, lease) = match self.reserve_exit_stream(&request, policy) {
518            Ok(reserved) => reserved,
519            Err(error) => {
520                return self
521                    .reject_exit_open(&request, OnionExitFailure::from_error(&error))
522                    .await;
523            }
524        };
525
526        let stream = match timeout(
527            Duration::from_secs(TCP_OPEN_TIMEOUT_SECS),
528            connect_exit_target(&target),
529        )
530        .await
531        {
532            Ok(Ok(stream)) => stream,
533            Ok(Err(failure)) => {
534                self.remove_exit_stream(request.key);
535                drop(lease);
536                return self.reject_exit_open(&request, failure).await;
537            }
538            Err(_) => {
539                self.remove_exit_stream(request.key);
540                drop(lease);
541                return self
542                    .reject_exit_open(&request, OnionExitFailure::ConnectTarget)
543                    .await;
544            }
545        };
546        if let Err(error) = self.accept_exit_open(&request).await {
547            self.remove_exit_stream(request.key);
548            drop(lease);
549            return Err(error);
550        }
551        let TcpExitOpen {
552            scope,
553            key,
554            circuit_id,
555            return_peer,
556            return_session_public_key,
557            client,
558            service,
559            ..
560        } = request;
561        spawn_exit_stream(ExitStreamTask {
562            runtime: self.clone(),
563            scope,
564            key,
565            circuit_id,
566            return_peer,
567            return_session_public_key,
568            client,
569            service,
570            stream,
571            rx,
572            lease,
573        });
574        Ok(())
575    }
576
577    async fn reject_exit_open(
578        &self,
579        request: &TcpExitOpen,
580        failure: OnionExitFailure,
581    ) -> Result<()> {
582        self.send_exit_backward(
583            request,
584            OnionBackwardSequence::FIRST,
585            OnionTcpPayload::Error(failure),
586        )
587        .await
588    }
589
590    async fn accept_exit_open(&self, request: &TcpExitOpen) -> Result<()> {
591        let sequence = self.next_backward_sequence(request.key)?;
592        self.send_exit_backward(request, sequence, OnionTcpPayload::Opened)
593            .await
594    }
595
596    async fn send_exit_backward(
597        &self,
598        request: &TcpExitOpen,
599        sequence: OnionBackwardSequence,
600        payload: OnionTcpPayload,
601    ) -> Result<()> {
602        // Resolve/connect results in the same quantum share one response deadline. The state and
603        // result algebra remain unchanged while remote clients lose byte-accurate resolver and
604        // target-connect timing.
605        tokio::time::sleep_until(open_response_deadline(request.opened_at, Instant::now())).await;
606        TcpBackwardRoute {
607            link_sender: &self.link_sender,
608            scope: &request.scope,
609            signer: &self.session_sk,
610            service: &request.service,
611            circuit_id: request.circuit_id,
612            return_peer: request.return_peer,
613            return_session_public_key: request.return_session_public_key,
614            client: request.client,
615        }
616        .send(sequence, payload)
617        .await
618    }
619
620    fn reserve_exit_stream(
621        &self,
622        request: &TcpExitOpen,
623        policy: &OnionExitPolicy,
624    ) -> Result<(mpsc::Receiver<TcpInbound>, OnionExitLease)> {
625        let (tx, rx) = mpsc::channel(32);
626        self.insert_exit_stream(
627            request.key,
628            request.service.clone(),
629            request.expected_forward_peer,
630            tx,
631        )?;
632        match self.admit_exit_stream(policy, request.circuit_id, request.return_peer, 0) {
633            Ok(lease) => Ok((rx, lease)),
634            Err(error) => {
635                self.remove_exit_stream(request.key);
636                Err(error)
637            }
638        }
639    }
640
641    fn insert_client_stream(
642        &self,
643        service: OnionServiceName,
644        expected_return_peer: Did,
645        expected_exit: OnionExitDescriptor,
646        return_id: OnionReturnId,
647        open_ack: oneshot::Sender<std::result::Result<(), OnionExitFailure>>,
648        tx: mpsc::Sender<TcpInbound>,
649    ) -> Result<TcpStreamKey> {
650        let mut streams = lock(&self.client_streams)?;
651        for _ in 0..16 {
652            let key = TcpStreamKey {
653                circuit_id: OnionCircuitId::random(),
654            };
655            match streams.entry(key) {
656                Entry::Vacant(entry) => {
657                    entry.insert(ClientStream {
658                        service,
659                        expected_return_peer,
660                        expected_exit,
661                        return_id,
662                        open_ack: Some(open_ack),
663                        backward_sequences: OnionSequenceWindow::default(),
664                        tx,
665                    });
666                    return Ok(key);
667                }
668                Entry::Occupied(_) => {}
669            }
670        }
671        Err(Error::OnionRouteError(
672            OnionRouteError::CircuitIdAllocationFailed,
673        ))
674    }
675
676    fn insert_exit_stream(
677        &self,
678        key: TcpStreamKey,
679        service: OnionServiceName,
680        expected_forward_peer: Did,
681        tx: mpsc::Sender<TcpInbound>,
682    ) -> Result<()> {
683        let mut streams = lock(&self.exit_streams)?;
684        match streams.entry(key) {
685            Entry::Vacant(entry) => {
686                entry.insert(ExitStream {
687                    service,
688                    expected_forward_peer,
689                    forward_sequences: OnionSequenceWindow::with_initial(
690                        OnionForwardSequence::FIRST.value(),
691                    ),
692                    next_backward_sequence: 0,
693                    tx,
694                });
695                Ok(())
696            }
697            Entry::Occupied(_) => Err(Error::OnionRouteError(OnionRouteError::DuplicateTcpOpen)),
698        }
699    }
700
701    fn send_client_inbound(&self, key: TcpStreamKey, from: Did, inbound: TcpInbound) -> Result<()> {
702        let tx = self.client_inbound_sender(key, from)?;
703        tx.try_send(inbound).map_err(|error| match error {
704            tokio::sync::mpsc::error::TrySendError::Full(_) => {
705                self.remove_client_stream(key);
706                Error::OnionRouteError(OnionRouteError::TcpStreamBackpressure)
707            }
708            tokio::sync::mpsc::error::TrySendError::Closed(_) => {
709                Error::OnionRouteError(OnionRouteError::TcpStreamClosed)
710            }
711        })
712    }
713
714    fn send_exit_inbound(
715        &self,
716        key: TcpStreamKey,
717        from: Did,
718        service: &OnionServiceName,
719        sequence: OnionForwardSequence,
720        inbound: TcpInbound,
721    ) -> Result<()> {
722        let tx = self.exit_inbound_sender(key, from, service, sequence)?;
723        tx.try_send(inbound).map_err(|error| match error {
724            tokio::sync::mpsc::error::TrySendError::Full(_) => {
725                self.remove_exit_stream(key);
726                Error::OnionRouteError(OnionRouteError::TcpStreamBackpressure)
727            }
728            tokio::sync::mpsc::error::TrySendError::Closed(_) => {
729                Error::OnionRouteError(OnionRouteError::TcpStreamClosed)
730            }
731        })
732    }
733
734    fn client_stream_service(&self, key: TcpStreamKey, from: Did) -> Result<OnionServiceName> {
735        let streams = lock(&self.client_streams)?;
736        let stream = authorize_client_stream(&streams, key, from)?;
737        Ok(stream.service.clone())
738    }
739
740    fn client_inbound_sender(
741        &self,
742        key: TcpStreamKey,
743        from: Did,
744    ) -> Result<mpsc::Sender<TcpInbound>> {
745        let streams = lock(&self.client_streams)?;
746        let stream = authorize_client_stream(&streams, key, from)?;
747        Ok(stream.tx.clone())
748    }
749
750    fn verify_client_payload(
751        &self,
752        key: TcpStreamKey,
753        from: Did,
754        payload: OnionAuthenticatedPayload,
755    ) -> Result<OnionCircuitPayload> {
756        let (service, expected_exit, return_id) = {
757            let streams = lock(&self.client_streams)?;
758            let stream = authorize_client_stream(&streams, key, from)?;
759            (
760                stream.service.clone(),
761                stream.expected_exit.clone(),
762                stream.return_id,
763            )
764        };
765        let verified = payload.into_verified_payload(return_id, &expected_exit)?;
766        if !verified.payload.is_service(&service) {
767            return Err(Error::OnionRouteError(
768                OnionRouteError::PayloadServiceMismatch {
769                    payload_service: verified.payload.service().to_string(),
770                    route_service: service.as_str().to_string(),
771                },
772            ));
773        }
774        self.consume_backward_sequence(key, from, verified.sequence)?;
775        Ok(verified.payload)
776    }
777
778    fn consume_backward_sequence(
779        &self,
780        key: TcpStreamKey,
781        from: Did,
782        sequence: OnionBackwardSequence,
783    ) -> Result<()> {
784        let mut streams = lock(&self.client_streams)?;
785        let stream = authorize_client_stream_mut(&mut streams, key, from)?;
786        match stream.backward_sequences.consume(sequence.value()) {
787            SequenceAdmission::Consumed => Ok(()),
788            SequenceAdmission::Duplicate => {
789                tracing::debug!(
790                    ?key,
791                    sequence = sequence.value(),
792                    "duplicate onion TCP backward sequence"
793                );
794                Err(Error::OnionRouteError(OnionRouteError::BackwardReplay))
795            }
796            SequenceAdmission::Stale => {
797                tracing::debug!(
798                    ?key,
799                    sequence = sequence.value(),
800                    "stale onion TCP backward sequence"
801                );
802                Err(Error::OnionRouteError(OnionRouteError::BackwardReplay))
803            }
804        }
805    }
806
807    fn complete_client_open(
808        &self,
809        key: TcpStreamKey,
810        from: Did,
811        result: std::result::Result<(), OnionExitFailure>,
812    ) -> Result<bool> {
813        let mut streams = lock(&self.client_streams)?;
814        let stream = authorize_client_stream_mut(&mut streams, key, from)?;
815        let Some(open_ack) = stream.open_ack.take() else {
816            return Ok(false);
817        };
818        let _ = open_ack.send(result);
819        Ok(true)
820    }
821
822    fn exit_inbound_sender(
823        &self,
824        key: TcpStreamKey,
825        from: Did,
826        service: &OnionServiceName,
827        sequence: OnionForwardSequence,
828    ) -> Result<mpsc::Sender<TcpInbound>> {
829        let mut streams = lock(&self.exit_streams)?;
830        let stream = authorize_exit_stream(&mut streams, key, from)?;
831        if &stream.service != service {
832            return Err(Error::OnionRouteError(
833                OnionRouteError::PayloadServiceMismatch {
834                    payload_service: service.as_str().to_string(),
835                    route_service: stream.service.as_str().to_string(),
836                },
837            ));
838        }
839        match stream.forward_sequences.consume(sequence.value()) {
840            SequenceAdmission::Consumed => {}
841            SequenceAdmission::Duplicate => {
842                tracing::debug!(
843                    ?key,
844                    sequence = sequence.value(),
845                    "duplicate onion TCP forward sequence"
846                );
847                return Err(Error::OnionRouteError(OnionRouteError::ForwardReplay));
848            }
849            SequenceAdmission::Stale => {
850                tracing::debug!(
851                    ?key,
852                    sequence = sequence.value(),
853                    "stale onion TCP forward sequence"
854                );
855                return Err(Error::OnionRouteError(OnionRouteError::ForwardReplay));
856            }
857        }
858        Ok(stream.tx.clone())
859    }
860
861    fn next_backward_sequence(&self, key: TcpStreamKey) -> Result<OnionBackwardSequence> {
862        let mut streams = lock(&self.exit_streams)?;
863        let stream = streams
864            .get_mut(&key)
865            .ok_or(Error::OnionRouteError(OnionRouteError::UnknownTcpStream))?;
866        let sequence = stream.next_backward_sequence;
867        stream.next_backward_sequence = sequence
868            .checked_add(1)
869            .ok_or(Error::OnionRouteError(OnionRouteError::SequenceExhausted))?;
870        Ok(OnionBackwardSequence::new(sequence))
871    }
872
873    fn remove_client_stream(&self, key: TcpStreamKey) {
874        if let Ok(mut streams) = self.client_streams.lock() {
875            streams.remove(&key);
876        }
877    }
878
879    fn remove_exit_stream(&self, key: TcpStreamKey) {
880        if let Ok(mut streams) = self.exit_streams.lock() {
881            streams.remove(&key);
882        }
883    }
884
885    fn admit_exit_stream(
886        &self,
887        policy: &OnionExitPolicy,
888        circuit_id: OnionCircuitId,
889        return_peer: Did,
890        bytes: u64,
891    ) -> Result<OnionExitLease> {
892        self.accounting
893            .admit(policy, circuit_id, return_peer, bytes)
894    }
895
896    fn record_exit_bytes(&self, policy: &OnionExitPolicy, bytes: u64) -> Result<()> {
897        self.accounting.record_bytes(policy, bytes)
898    }
899}
900
901#[derive(Clone, Copy, Debug, Eq, Hash, PartialEq)]
902struct TcpStreamKey {
903    circuit_id: OnionCircuitId,
904}
905
906struct TcpExitOpen {
907    scope: Scope,
908    opened_at: Instant,
909    key: TcpStreamKey,
910    circuit_id: OnionCircuitId,
911    return_peer: Did,
912    return_session_public_key: PublicKey<33>,
913    client: OnionClientReturn,
914    expected_forward_peer: Did,
915    service: OnionServiceName,
916    target: String,
917}
918
919// Invariant: each sequence in `backward_sequences` has already produced at most one
920// `TcpInbound` event for this client stream.
921// Preservation: `verify_client_payload` verifies the exit proof and consumes the monotonic
922// sequence before decoding the TCP payload; duplicate/stale sequences fail before bytes reach the
923// stream.
924// Invariant: `service` is the canonical route service used for every client-to-exit payload on this
925// stream.
926// Preservation: `verify_client_payload` rejects signed backward payloads whose service differs
927// from this stream service before bytes reach the stream.
928struct ClientStream {
929    service: OnionServiceName,
930    expected_return_peer: Did,
931    expected_exit: OnionExitDescriptor,
932    return_id: OnionReturnId,
933    open_ack: Option<oneshot::Sender<std::result::Result<(), OnionExitFailure>>>,
934    backward_sequences: OnionSequenceWindow,
935    tx: mpsc::Sender<TcpInbound>,
936}
937
938// Invariant: `service` is the canonical service accepted by the Open payload that created this exit
939// stream.
940// Preservation: `exit_inbound_sender` rejects later payloads on the same circuit when their service
941// differs from this stream service.
942struct ExitStream {
943    service: OnionServiceName,
944    expected_forward_peer: Did,
945    forward_sequences: OnionSequenceWindow,
946    next_backward_sequence: u64,
947    tx: mpsc::Sender<TcpInbound>,
948}
949
950fn authorize_client_stream(
951    streams: &HashMap<TcpStreamKey, ClientStream>,
952    key: TcpStreamKey,
953    actual: Did,
954) -> Result<&ClientStream> {
955    let stream = streams
956        .get(&key)
957        .ok_or(Error::OnionRouteError(OnionRouteError::UnknownTcpStream))?;
958    if stream.expected_return_peer != actual {
959        return Err(Error::OnionRouteError(
960            OnionRouteError::UnexpectedTcpReturnPeer {
961                expected: stream.expected_return_peer,
962                actual,
963            },
964        ));
965    }
966    Ok(stream)
967}
968
969fn authorize_client_stream_mut(
970    streams: &mut HashMap<TcpStreamKey, ClientStream>,
971    key: TcpStreamKey,
972    actual: Did,
973) -> Result<&mut ClientStream> {
974    let stream = streams
975        .get_mut(&key)
976        .ok_or(Error::OnionRouteError(OnionRouteError::UnknownTcpStream))?;
977    if stream.expected_return_peer != actual {
978        return Err(Error::OnionRouteError(
979            OnionRouteError::UnexpectedTcpReturnPeer {
980                expected: stream.expected_return_peer,
981                actual,
982            },
983        ));
984    }
985    Ok(stream)
986}
987
988fn authorize_exit_stream(
989    streams: &mut HashMap<TcpStreamKey, ExitStream>,
990    key: TcpStreamKey,
991    actual: Did,
992) -> Result<&mut ExitStream> {
993    let stream = streams
994        .get_mut(&key)
995        .ok_or(Error::OnionRouteError(OnionRouteError::UnknownTcpStream))?;
996    if stream.expected_forward_peer != actual {
997        return Err(Error::OnionRouteError(
998            OnionRouteError::UnexpectedTcpForwardPeer {
999                expected: stream.expected_forward_peer,
1000                actual,
1001            },
1002        ));
1003    }
1004    Ok(stream)
1005}
1006
1007#[cfg(test)]
1008mod tests;