arti/subcommands/
proxy.rs1use std::sync::Arc;
4
5use anyhow::{Context, Result};
6use cfg_if::cfg_if;
7use clap::ArgMatches;
8#[allow(unused)]
9use tor_config_path::CfgPathResolver;
10use tracing::{info, instrument, warn};
11
12use arti_client::TorClientConfig;
13use tor_config::{ConfigurationSources, Listen};
14use tor_rtcompat::ToplevelRuntime;
15
16#[cfg(feature = "dns-proxy")]
17use crate::dns;
18use crate::{
19 ArtiConfig, TorClient, exit, process,
20 proxy::{self, ListenProtocols, port_info},
21 reload_cfg,
22};
23
24#[cfg(feature = "rpc")]
25use crate::rpc;
26
27#[cfg(feature = "onion-service-service")]
28use crate::onion_proxy;
29
30type PinnedFuture<T> = std::pin::Pin<Box<dyn futures::Future<Output = T>>>;
32
33#[instrument(skip_all, level = "trace")]
35pub(crate) fn run<R: ToplevelRuntime>(
36 runtime: R,
37 proxy_matches: &ArgMatches,
38 cfg_sources: ConfigurationSources,
39 loaded_cfg: tor_config::ConfigurationTree,
40 config: ArtiConfig,
41 client_config: TorClientConfig,
42) -> Result<()> {
43 let socks_listen = match proxy_matches.get_one::<u16>("socks-port") {
48 Some(p) => Listen::new_localhost(*p),
49 None => config.proxy().socks_listen.clone(),
50 };
51
52 let dns_listen = match proxy_matches.get_one::<u16>("dns-port") {
54 Some(p) => Listen::new_localhost(*p),
55 None => config.proxy().dns_listen.clone(),
56 };
57
58 if !socks_listen.is_empty() {
59 info!(
60 "Starting Arti {} in proxy mode on {} ...",
61 env!("CARGO_PKG_VERSION"),
62 socks_listen
63 );
64 }
65
66 if let Some(listen) = {
67 config
69 .metrics
70 .prometheus
71 .listen
72 .single_address_legacy()
73 .context("can only listen on a single address for Prometheus metrics")?
74 } {
75 cfg_if! {
76 if #[cfg(feature = "metrics")] {
77 metrics_exporter_prometheus::PrometheusBuilder::new()
78 .with_http_listener(listen)
79 .install()
80 .with_context(|| format!(
81 "set up Prometheus metrics exporter on {listen}"
82 ))?;
83 info!("Arti Prometheus metrics export scraper endpoint http://{listen}");
84 } else {
85 return Err(anyhow::anyhow!(
86 "`metrics.prometheus.listen` config set but `metrics` cargo feature compiled out in `arti` crate"
87 ));
88 }
89 }
90 }
91
92 #[cfg(not(all(target_arch = "wasm32", target_os = "unknown")))]
93 process::use_max_file_limit(&config);
94
95 let rt_copy = runtime.clone();
96 rt_copy.block_on(run_proxy(
97 runtime,
98 socks_listen,
99 dns_listen,
100 config.proxy().protocols(),
101 cfg_sources,
102 loaded_cfg,
103 config,
104 client_config,
105 ))?;
106
107 Ok(())
108}
109
110#[cfg_attr(feature = "experimental-api", visibility::make(pub))]
116#[cfg_attr(docsrs, doc(cfg(feature = "experimental-api")))]
117#[instrument(skip_all, level = "trace")]
118#[expect(clippy::too_many_arguments)]
119async fn run_proxy<R: ToplevelRuntime>(
120 runtime: R,
121 socks_listen: Listen,
122 dns_listen: Listen,
123 protocols: ListenProtocols,
124 config_sources: ConfigurationSources,
127 loaded_cfg: tor_config::ConfigurationTree,
128 arti_config: ArtiConfig,
129 client_config: TorClientConfig,
130) -> Result<()> {
131 use arti_client::BootstrapBehavior;
134 use futures::FutureExt;
135
136 let fs_mistrust = client_config.fs_mistrust().clone();
138 let path_resolver: CfgPathResolver = AsRef::<CfgPathResolver>::as_ref(&client_config).clone();
139
140 let defer_bootstrap = arti_config.application().defer_bootstrap;
141
142 let bootstrap_behavior = match defer_bootstrap {
143 true => BootstrapBehavior::Manual,
144 false => BootstrapBehavior::OnDemand,
145 };
146
147 let (cfg_mgr, cfg_watcher_task) = reload_cfg::CfgMgr::new(
148 runtime.clone(),
149 config_sources,
150 #[cfg(feature = "rpc")]
151 loaded_cfg,
152 &arti_config,
153 vec![],
154 )?;
155
156 let client_builder = TorClient::with_runtime(runtime.clone())
157 .config(client_config)
158 .bootstrap_behavior(bootstrap_behavior);
159 let client = client_builder.create_unbootstrapped_async().await?;
160
161 let launchable_client = Arc::new(reload_cfg::LaunchableTorClient::new(
162 Arc::clone(&client),
163 arti_config.application(),
164 ));
165
166 #[allow(unused_mut)]
167 let mut reconfigurable_modules: Vec<Arc<dyn reload_cfg::ReconfigurableModule>> = vec![
168 Arc::clone(&launchable_client) as _,
169 Arc::new(reload_cfg::Application::new(arti_config.clone())),
170 ];
171
172 cfg_if::cfg_if! {
173 if #[cfg(feature = "onion-service-service")] {
174 let have_onion_svc = if defer_bootstrap {
175 let onion_services = onion_proxy::ProxySet::new_deferred(Arc::clone(&client));
176 reconfigurable_modules.push(Arc::new(onion_services));
177 arti_config.onion_services.values().any(|c| *c.svc_cfg.enabled())
178 } else {
179 let onion_services =
180 onion_proxy::ProxySet::launch_new(Arc::clone(&client), arti_config.onion_services.clone())?;
181 let have_onion_svc = !onion_services.is_empty();
182 reconfigurable_modules.push(Arc::new(onion_services));
183 have_onion_svc
184 };
185 } else {
186 let have_onion_svc = false;
187 }
188 };
189
190 reconfigurable_modules
197 .iter()
198 .try_for_each(|m| cfg_watcher_task.add_module(m))?;
199
200 cfg_watcher_task.launch()?;
201
202 cfg_if::cfg_if! {
203 if #[cfg(feature = "rpc")] {
204 let rpc_data = rpc::launch_rpc_mgr(
205 &runtime,
206 &arti_config.rpc,
207 &path_resolver,
208 &fs_mistrust,
209 client.clone(),
210 launchable_client.clone(),
211 cfg_mgr.clone(),
212 )
213 .await?;
214 let (rpc_mgr, mut rpc_state_sender) = rpc_data
215 .map(|d| (d.rpc_mgr, d.rpc_state_sender))
216 .unzip();
217 } else {
218 let rpc_mgr = None;
219 }
220 }
221
222 let mut listen_options = tor_rtcompat::TcpListenOptions::builder();
224 listen_options
225 .common()
226 .send_buffer_size(Some(arti_config.proxy().socket_send_buf_size.as_usize()))
227 .recv_buffer_size(Some(arti_config.proxy().socket_recv_buf_size.as_usize()));
228 let listen_options = listen_options.build()?;
229
230 let mut proxy: Vec<PinnedFuture<Result<()>>> = Vec::new();
231 let mut ports = Vec::new();
232 if !socks_listen.is_empty() {
233 let runtime = runtime.clone();
234 let client = client.isolated_client();
235 let socks_listen = socks_listen.clone();
236 let listener_type = protocols.to_string();
237
238 let stream_proxy = proxy::bind_proxy(
239 runtime,
240 client,
241 socks_listen,
242 listen_options,
243 protocols,
244 rpc_mgr,
245 )
246 .await
247 .with_context(|| format!("Unable to launch {listener_type} proxy"))?;
248 let port_info = stream_proxy.port_info()?;
249
250 ports.extend(port_info);
251
252 let failure_message = format!("{listener_type} proxy died unexpectedly");
253 let proxy_future = stream_proxy
254 .run_proxy()
255 .map(|future_result| future_result.context(failure_message));
256 proxy.push(Box::pin(proxy_future));
257 }
258
259 #[cfg(feature = "dns-proxy")]
260 if !dns_listen.is_empty() {
261 let runtime = runtime.clone();
262 let client = client.isolated_client();
263 let dns_proxy = dns::bind_dns_resolver(runtime, client, dns_listen)
264 .await
265 .context("Unable to launch DNS proxy")?;
266 ports.extend(dns_proxy.port_info().context("Unable to find DNS ports")?);
267 let proxy_future = dns_proxy
268 .run_dns_proxy()
269 .map(|future_result| future_result.context("DNS proxy died unexpectedly"));
270 proxy.push(Box::pin(proxy_future));
271 }
272
273 #[cfg(not(feature = "dns-proxy"))]
274 if !dns_listen.is_empty() {
275 warn!(
276 "Tried to specify a DNS proxy address, but Arti was built without dns-proxy support."
277 );
278 return Ok(());
279 }
280
281 if proxy.is_empty() {
282 if !have_onion_svc {
283 warn!(
285 "No proxy address set; \
286 specify -p PORT (to override `socks_listen`) \
287 or -d PORT (to override `dns_listen`). \
288 Alternatively, use the `socks_listen` or `dns_listen` configuration options."
289 );
290 return Ok(());
291 } else {
292 proxy.push(Box::pin(futures::future::pending()));
295 }
296 }
297
298 cfg_if::cfg_if! {
299 if #[cfg(feature="rpc")] {
300 if let Some(rpc_state_sender) = &mut rpc_state_sender {
301 rpc_state_sender.set_stream_listeners(&ports[..]);
302 }
303 }
304 }
305
306 {
307 let port_info = port_info::PortInfo { ports };
308 let port_info_file = arti_config
309 .storage()
310 .port_info_file
311 .path(&path_resolver)
312 .context("Can't find path for port_info_file")?;
313 if port_info_file.to_str() != Some("") {
314 port_info.write_to_file(&fs_mistrust, &port_info_file)?;
315 }
316 }
317
318 let proxy = futures::future::select_all(proxy).map(|(finished, _index, _others)| finished);
319 futures::select!(
320 r = exit::wait_for_ctrl_c().fuse()
321 => r.context("waiting for termination signal"),
322 r = proxy.fuse()
323 => r,
324 r = async {
325 if defer_bootstrap {
326 info!("Bootstrapping deferred.");
327 } else {
328 client.bootstrap().await?;
329 if !socks_listen.is_empty() {
330 info!("Sufficiently bootstrapped; proxy now functional.");
331 } else {
332 info!("Sufficiently bootstrapped.");
333 }
334 }
335 futures::future::pending::<Result<()>>().await
336 }.fuse()
337 => r.context("bootstrap"),
338 )?;
339
340 drop(reconfigurable_modules);
344 drop(cfg_mgr);
345
346 Ok(())
347}