Skip to main content

cuprate_p2p_transport/
arti.rs

1//! Arti Transport
2//!
3//! This module defines a transport method for the `Tor` network zone using the `arti_client` library.
4//!
5
6//---------------------------------------------------------------------------------------------------- Imports
7
8use std::{
9    io::{self, ErrorKind},
10    pin::Pin,
11    sync::Arc,
12    task::{Context, Poll},
13};
14
15use arti_client::{DataReader, DataWriter, TorClient, TorClientConfig};
16use async_trait::async_trait;
17use futures::{Stream, StreamExt};
18use tokio_util::codec::{FramedRead, FramedWrite};
19use tor_cell::relaycell::msg::Connected;
20use tor_config_path::CfgPathResolver;
21use tor_hsservice::{handle_rend_requests, OnionService, RunningOnionService};
22use tor_proto::stream::IncomingStreamRequest;
23use tor_rtcompat::PreferredRuntime;
24
25use cuprate_p2p_core::{ClearNet, NetworkZone, Tor, Transport};
26use cuprate_wire::MoneroWireCodec;
27
28use crate::DisabledListener;
29
30//---------------------------------------------------------------------------------------------------- Configuration
31
32#[derive(Clone)]
33pub struct ArtiClientConfig {
34    /// Arti bootstrapped client
35    pub client: Arc<TorClient<PreferredRuntime>>,
36}
37
38pub struct ArtiServerConfig {
39    /// Arti onion service
40    pub onion_svc: OnionService,
41    /// Listening port
42    pub port: u16,
43
44    // Mandatory resources for launching the onion service
45    client: Arc<TorClient<PreferredRuntime>>,
46    path_resolver: Arc<CfgPathResolver>,
47}
48
49impl ArtiServerConfig {
50    pub fn new(
51        onion_svc: OnionService,
52        port: u16,
53        client: &Arc<TorClient<PreferredRuntime>>,
54        config: &TorClientConfig,
55    ) -> Self {
56        let path_resolver: &CfgPathResolver = config.as_ref();
57
58        Self {
59            onion_svc,
60            port,
61            client: Arc::clone(client),
62            path_resolver: Arc::new(path_resolver.clone()),
63        }
64    }
65}
66
67//---------------------------------------------------------------------------------------------------- Transport
68
69type PinnedStream<I> = Pin<Box<dyn Stream<Item = I> + Send>>;
70
71/// An onion service listening for incoming peer connections.
72pub struct OnionListener {
73    /// A handle to the onion service instance.
74    _onion_svc: Arc<RunningOnionService>,
75    /// A modified stream that produce a data stream and sink from rendez-vous requests.
76    listener: PinnedStream<Result<(DataReader, DataWriter), io::Error>>,
77}
78
79impl Stream for OnionListener {
80    type Item = Result<
81        (
82            Option<<Tor as NetworkZone>::Addr>,
83            FramedRead<DataReader, MoneroWireCodec>,
84            FramedWrite<DataWriter, MoneroWireCodec>,
85        ),
86        io::Error,
87    >;
88
89    fn poll_next(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Option<Self::Item>> {
90        match self.listener.poll_next_unpin(cx) {
91            Poll::Pending => Poll::Pending,
92            Poll::Ready(req) => Poll::Ready(req.map(|r| {
93                r.map(|(stream, sink)| {
94                    (
95                        None, // Inbound is anonymous
96                        FramedRead::new(stream, MoneroWireCodec::default()),
97                        FramedWrite::new(sink, MoneroWireCodec::default()),
98                    )
99                })
100            })),
101        }
102    }
103}
104
105#[derive(Clone, Copy)]
106pub struct Arti;
107
108#[async_trait]
109impl Transport<Tor> for Arti {
110    type ClientConfig = ArtiClientConfig;
111    type ServerConfig = ArtiServerConfig;
112
113    type Stream = FramedRead<DataReader, MoneroWireCodec>;
114    type Sink = FramedWrite<DataWriter, MoneroWireCodec>;
115    type Listener = OnionListener;
116
117    async fn connect_to_peer(
118        addr: <Tor as NetworkZone>::Addr,
119        config: &Self::ClientConfig,
120    ) -> Result<(Self::Stream, Self::Sink), io::Error> {
121        config
122            .client
123            .connect((addr.addr_string(), addr.port()))
124            .await
125            .map_err(|e| io::Error::new(ErrorKind::ConnectionAborted, e.to_string()))
126            .map(|stream| {
127                let (stream, sink) = stream.split();
128                (
129                    FramedRead::new(stream, MoneroWireCodec::default()),
130                    FramedWrite::new(sink, MoneroWireCodec::default()),
131                )
132            })
133    }
134
135    async fn incoming_connection_listener(
136        config: Self::ServerConfig,
137    ) -> Result<Self::Listener, io::Error> {
138        let not_running =
139            |e: arti_client::Error| io::Error::new(ErrorKind::NotConnected, e.to_string());
140
141        // Launch onion service
142        let (svc, rdv_stream) = config
143            .onion_svc
144            .launch(
145                config.client.runtime().clone(),
146                config.client.dirmgr().map_err(not_running)?,
147                config.client.hs_circ_pool().map_err(not_running)?,
148                config.path_resolver,
149            )
150            .unwrap()
151            .expect("the onion service is never disabled in its config");
152
153        // Accept all rendez-vous and await correct stream request
154        #[expect(clippy::wildcard_enum_match_arm)]
155        let req_stream = handle_rend_requests(rdv_stream).then(move |sreq| async move {
156            match sreq.request() {
157                // As specified in: <https://spec.torproject.org/rend-spec/managing-streams.html>
158                //
159                // A client that wishes to open a data stream with us needs to send a BEGIN message with an empty address
160                // and no flags. We additionally filter requests to the correct port configured and advertised on P2P.
161                IncomingStreamRequest::Begin(r)
162                    if r.port() == config.port && r.addr().is_empty() && r.flags().is_empty() =>
163                {
164                    let stream = sreq
165                        .accept(Connected::new_empty())
166                        .await
167                        .map_err(|e| io::Error::new(ErrorKind::BrokenPipe, e.to_string()))?;
168
169                    Ok(stream.split())
170                }
171                _ => {
172                    sreq.shutdown_circuit()
173                        .expect("Should never panic, unless programming error from arti's end.");
174
175                    Err(io::Error::other("Received invalid command"))
176                }
177            }
178        });
179
180        Ok(OnionListener {
181            _onion_svc: svc,
182            listener: Box::pin(req_stream),
183        })
184    }
185}
186
187#[async_trait]
188impl Transport<ClearNet> for Arti {
189    type ClientConfig = ArtiClientConfig;
190    type ServerConfig = ();
191
192    type Stream = FramedRead<DataReader, MoneroWireCodec>;
193    type Sink = FramedWrite<DataWriter, MoneroWireCodec>;
194    type Listener = DisabledListener<ClearNet, DataReader, DataWriter>;
195
196    async fn connect_to_peer(
197        addr: <ClearNet as NetworkZone>::Addr,
198        config: &Self::ClientConfig,
199    ) -> Result<(Self::Stream, Self::Sink), io::Error> {
200        config
201            .client
202            .connect(addr.to_string())
203            .await
204            .map_err(|e| io::Error::new(ErrorKind::ConnectionAborted, e.to_string()))
205            .map(|stream| {
206                let (stream, sink) = stream.split();
207                (
208                    FramedRead::new(stream, MoneroWireCodec::default()),
209                    FramedWrite::new(sink, MoneroWireCodec::default()),
210                )
211            })
212    }
213
214    async fn incoming_connection_listener(
215        _config: Self::ServerConfig,
216    ) -> Result<Self::Listener, io::Error> {
217        panic!("In anonymized clearnet mode, inbound is disabled!");
218    }
219}