Skip to main content

arti/subcommands/
proxy.rs

1//! The `proxy` subcommand.
2
3use 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
30/// Shorthand for a boxed and pinned Future.
31type PinnedFuture<T> = std::pin::Pin<Box<dyn futures::Future<Output = T>>>;
32
33/// Run the `proxy` subcommand.
34#[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    // Override configured listen addresses from the command line.
44    // This implies listening on localhost ports.
45
46    // TODO: Parse a string rather than calling new_localhost.
47    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    // TODO: Parse a string rather than calling new_localhost.
53    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        // https://github.com/metrics-rs/metrics/issues/567
68        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/// Run the main loop of the proxy.
111///
112/// # Panics
113///
114/// Currently, might panic if things go badly enough wrong
115#[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    // TODO RPC Config: We are passing numerous types here in order to construct a CfgMgr.
125    // Can we instead construct one earlier?  Or at least put them in a struct?
126    config_sources: ConfigurationSources,
127    loaded_cfg: tor_config::ConfigurationTree,
128    arti_config: ArtiConfig,
129    client_config: TorClientConfig,
130) -> Result<()> {
131    // Using OnDemand arranges that, while we are bootstrapping, incoming connections wait
132    // for bootstrap to complete, rather than getting errors.
133    use arti_client::BootstrapBehavior;
134    use futures::FutureExt;
135
136    // TODO: We may instead want to provide a way to get these items out of TorClient.
137    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    // The add_module function will use references here
191    // to prevent the task spawned by watch_for_config_changes from
192    // keeping these modules alive after this function exits.
193    //
194    // NOTE: reconfigurable_modules stores the only strong references to these modules,
195    // so we must keep that variable alive until the end of the function
196    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    // The options that we'll use for our listening proxy sockets.
223    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            // TODO: rename "socks_listen" to "proxy_listen", preserving compat, once http-connect is stable.
284            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            // Push a dummy future to appease future::select_all,
293            // which expects a non-empty list
294            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    // The modules and CfgMgr can be dropped now, because we are exiting.
341    // (We drop them explicitly to make sure that they were not dropped
342    // accidentally before.)
343    drop(reconfigurable_modules);
344    drop(cfg_mgr);
345
346    Ok(())
347}