1use 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#[derive(Clone)]
33pub struct ArtiClientConfig {
34 pub client: Arc<TorClient<PreferredRuntime>>,
36}
37
38pub struct ArtiServerConfig {
39 pub onion_svc: OnionService,
41 pub port: u16,
43
44 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
67type PinnedStream<I> = Pin<Box<dyn Stream<Item = I> + Send>>;
70
71pub struct OnionListener {
73 _onion_svc: Arc<RunningOnionService>,
75 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, 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 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 #[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 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}