Skip to main content

arti/
rpc.rs

1//! Experimental RPC support.
2
3use anyhow::Result;
4use arti_rpcserver::RpcMgr;
5use derive_deftly::Deftly;
6use fs_mistrust::Mistrust;
7use futures::{AsyncReadExt, stream::StreamExt};
8use session::ArtiRpcSession;
9use std::collections::BTreeMap;
10use std::{io::Result as IoResult, sync::Arc};
11use tor_config::derive::prelude::*;
12use tor_config_path::CfgPathResolver;
13use tracing::{debug, info};
14
15use arti_client::TorClient;
16use tor_rtcompat::{NetStreamListener as _, Runtime, SpawnExt, general};
17
18semipublic_mod! { pub(crate) mod configuration; }
19pub(crate) mod conntarget;
20pub(crate) mod listener;
21mod proxyinfo;
22mod session;
23mod superuser;
24
25pub(crate) use configuration::{ConfigSettings, ConfigValue};
26use listener::RpcListenerSetConfig;
27pub(crate) use session::{RpcStateSender, RpcVisibleArtiState};
28
29use crate::reload_cfg::{CfgMgr, LaunchableTorClient};
30use crate::rpc::superuser::RpcSuperuser;
31
32/// Configuration for Arti's RPC subsystem.
33///
34/// You cannot change this section on a running Arti client.
35#[derive(Debug, Clone, Deftly, Eq, PartialEq)]
36#[derive_deftly(TorConfig)]
37#[cfg_attr(feature = "experimental-api", visibility::make(pub))]
38#[cfg_attr(feature = "experimental-api", deftly(tor_config(vis = pub)))]
39pub(crate) struct RpcConfig {
40    /// If true, then the RPC subsystem is enabled and will listen for connections.
41    #[deftly(tor_config(default = false))] // TODO RPC make this true once we are stable.
42    enable: bool,
43
44    /// A set of named locations in which to find connect files.
45    #[deftly(tor_config(map, default = listener::listener_map_defaults()))]
46    listen: BTreeMap<String, RpcListenerSetConfig>,
47
48    /// A list of default connect points to bind
49    /// if no enabled connect points are found under `listen`.
50    #[deftly(tor_config(list(element(clone)), default = listen_defaults_defaults()))]
51    listen_default: Vec<String>,
52}
53
54/// Return default values for `RpcConfig.listen_default`
55fn listen_defaults_defaults() -> Vec<String> {
56    vec![tor_rpc_connect::USER_DEFAULT_CONNECT_POINT.to_string()]
57}
58
59/// Information about an incoming connection.
60///
61/// Yielded in a stream from our RPC listeners.
62type IncomingConn = (
63    general::Stream,
64    general::SocketAddr,
65    Arc<listener::RpcConnInfo>,
66);
67
68/// Bind to all configured RPC listeners in `cfg`.
69///
70/// On success, return a stream of `IncomingConn`.
71async fn launch_all_listeners<R: Runtime>(
72    runtime: &R,
73    cfg: &RpcConfig,
74    resolver: &CfgPathResolver,
75    mistrust: &Mistrust,
76) -> anyhow::Result<(
77    impl futures::Stream<Item = IoResult<IncomingConn>> + Unpin + use<R>,
78    Vec<tor_rpc_connect::server::ListenerGuard>,
79)> {
80    let mut listeners = Vec::new();
81    let mut guards = Vec::new();
82    for (name, listener_cfg) in cfg.listen.iter() {
83        for (lis, info, guard) in listener_cfg
84            .bind(runtime, name.as_str(), resolver, mistrust)
85            .await?
86        {
87            // (Note that `bind` only returns enabled listeners, so we don't need to check here.
88            debug!(
89                "Listening at {} for {}",
90                lis.local_addr()
91                    .expect("general::listener without address?")
92                    .display_lossy(),
93                info.name,
94            );
95            listeners.push((lis, info));
96            guards.push(guard);
97        }
98    }
99    if listeners.is_empty() {
100        for (idx, connpt) in cfg.listen_default.iter().enumerate() {
101            let display_index = idx + 1; // One-indexed values are more human-readable.
102            let (lis, info, guard) =
103                listener::bind_string(connpt, display_index, runtime, resolver, mistrust).await?;
104            debug!(
105                "Listening at {} for {}",
106                lis.local_addr()
107                    .expect("general::listener without address?")
108                    .display_lossy(),
109                info.name,
110            );
111            listeners.push((lis, info));
112            guards.push(guard);
113        }
114    }
115    if listeners.is_empty() {
116        info!("No RPC listeners configured.");
117    }
118
119    let streams = listeners.into_iter().map(|(listener, info)| {
120        listener
121            .incoming()
122            .map(move |accept_result| match accept_result {
123                Ok((netstream, addr)) => Ok((netstream, addr, Arc::clone(&info))),
124                Err(e) => Err(e),
125            })
126    });
127
128    Ok((futures::stream::select_all(streams), guards))
129}
130
131/// Create an RPC manager, bind to connect points, and open a listener task to accept incoming
132/// RPC connections.
133pub(crate) async fn launch_rpc_mgr<R: Runtime>(
134    runtime: &R,
135    cfg: &RpcConfig,
136    resolver: &CfgPathResolver,
137    mistrust: &Mistrust,
138    client: Arc<TorClient<R>>,
139    launchable: Arc<LaunchableTorClient<R>>,
140    cfg_mgr: Arc<CfgMgr<R>>,
141) -> Result<Option<RpcProxySupport>> {
142    if !cfg.enable {
143        return Ok(None);
144    }
145    let (rpc_state, rpc_state_sender) = RpcVisibleArtiState::new();
146
147    let rpc_mgr = RpcMgr::new()?;
148    // Register methods. Needed since TorClient is generic.
149    //
150    // TODO: If we accumulate a large number of generics like this, we should do this elsewhere.
151    rpc_mgr.register_rpc_methods(TorClient::<R>::rpc_methods());
152    rpc_mgr.register_rpc_methods(arti_rpcserver::rpc_methods::<R>());
153    rpc_mgr.register_rpc_methods(RpcSuperuser::<R>::rpc_methods());
154
155    let rt_clone = runtime.clone();
156    let rpc_mgr_clone = rpc_mgr.clone();
157
158    let (incoming, guards) = launch_all_listeners(runtime, cfg, resolver, mistrust).await?;
159
160    // TODO: Using spawn in this way makes it hard to report whether we
161    // succeeded or not. This is something we should fix when we refactor
162    // our service-launching code.
163    runtime.spawn(async move {
164        let result = run_rpc_listener(
165            rt_clone,
166            incoming,
167            rpc_mgr_clone,
168            client,
169            launchable,
170            cfg_mgr,
171            rpc_state,
172        )
173        .await;
174        if let Err(e) = result {
175            tracing::warn!("RPC manager quit with an error: {}", e);
176        }
177        drop(guards);
178    })?;
179    Ok(Some(RpcProxySupport {
180        rpc_mgr,
181        rpc_state_sender,
182    }))
183}
184
185/// Backend function to implement an RPC listener: runs in a loop.
186async fn run_rpc_listener<R: Runtime>(
187    runtime: R,
188    mut incoming: impl futures::Stream<Item = IoResult<IncomingConn>> + Unpin,
189    rpc_mgr: Arc<RpcMgr>,
190    client: Arc<TorClient<R>>,
191    launchable: Arc<LaunchableTorClient<R>>,
192    cfg_mgr: Arc<CfgMgr<R>>,
193    rpc_state: Arc<RpcVisibleArtiState>,
194) -> Result<()> {
195    while let Some((stream, _addr, info)) = incoming.next().await.transpose()? {
196        debug!("Received incoming RPC connection from {}", &info.name);
197
198        let client_clone = client.clone();
199        let rpc_state_clone = rpc_state.clone();
200        let launchable = launchable.clone();
201        let cfg_mgr_clone = cfg_mgr.clone();
202        let connection = rpc_mgr.new_connection(info.auth.clone(), move |auth| {
203            ArtiRpcSession::new(
204                auth,
205                &client_clone,
206                &launchable,
207                &rpc_state_clone,
208                &cfg_mgr_clone,
209                &info,
210            ) as _
211        });
212        let (input, output) = stream.split();
213
214        runtime.spawn(async {
215            let result = connection.run(input, output).await;
216            if let Err(e) = result {
217                tracing::warn!("RPC session ended with an error: {}", e);
218            }
219        })?;
220    }
221    Ok(())
222}
223
224/// Information passed to a proxy or similar stream provider when running with RPC support.
225pub(crate) struct RpcProxySupport {
226    /// An RPC manager to use for looking up objects as possible stream targets.
227    pub(crate) rpc_mgr: Arc<arti_rpcserver::RpcMgr>,
228    /// An RPCStateSender to use for registering the list of known proxy ports.
229    pub(crate) rpc_state_sender: RpcStateSender,
230}
231
232#[cfg(test)]
233mod test {
234    // @@ begin test lint list maintained by maint/add_warning @@
235    #![allow(clippy::bool_assert_comparison)]
236    #![allow(clippy::clone_on_copy)]
237    #![allow(clippy::dbg_macro)]
238    #![allow(clippy::mixed_attributes_style)]
239    #![allow(clippy::print_stderr)]
240    #![allow(clippy::print_stdout)]
241    #![allow(clippy::single_char_pattern)]
242    #![allow(clippy::unwrap_used)]
243    #![allow(clippy::unchecked_time_subtraction)]
244    #![allow(clippy::useless_vec)]
245    #![allow(clippy::needless_pass_by_value)]
246    #![allow(clippy::string_slice)] // See arti#2571
247    //! <!-- @@ end test lint list maintained by maint/add_warning @@ -->
248
249    use listener::{ConnectPointOptionsBuilder, RpcListenerSetConfigBuilder};
250    use tor_config_path::CfgPath;
251    use tor_rpc_connect::ParsedConnectPoint;
252
253    use super::*;
254
255    #[test]
256    fn rpc_method_names() {
257        // We run this from a nice high level module, to ensure that as many method names as
258        // possible will be in-scope.
259        let problems = tor_rpcbase::check_method_names([]);
260
261        for (m, err) in &problems {
262            eprintln!("Bad method name {m:?}: {err}");
263        }
264        assert!(problems.is_empty());
265    }
266
267    #[test]
268    fn parse_listener_defaults() {
269        for string in listen_defaults_defaults() {
270            let _parsed: ParsedConnectPoint = string.parse().unwrap();
271        }
272    }
273
274    #[test]
275    fn parsing_and_building() {
276        fn build(s: &str) -> Result<RpcConfig, anyhow::Error> {
277            let b: RpcConfigBuilder = toml::from_str(s)?;
278            Ok(b.build()?)
279        }
280
281        let mut user_defaults_builder = RpcListenerSetConfigBuilder::default();
282        user_defaults_builder.listener_options().enable(true);
283        user_defaults_builder.dir(CfgPath::new("${ARTI_LOCAL_DATA}/rpc/connect.d".to_string()));
284        let mut system_defaults_builder = RpcListenerSetConfigBuilder::default();
285        system_defaults_builder.listener_options().enable(false);
286        system_defaults_builder.dir(CfgPath::new("/etc/arti-rpc/connect.d".to_string()));
287
288        // Make sure that an empty configuration gets us the defaults.
289        let defaults = build("").unwrap();
290        assert_eq!(
291            defaults,
292            RpcConfig {
293                enable: false,
294                listen: vec![
295                    (
296                        "user-default".to_string(),
297                        user_defaults_builder.build().unwrap()
298                    ),
299                    (
300                        "system-default".to_string(),
301                        system_defaults_builder.build().unwrap()
302                    ),
303                ]
304                .into_iter()
305                .collect(),
306                listen_default: listen_defaults_defaults()
307            }
308        );
309
310        // Make sure that overriding specific options works as expected.
311        let altered = build(
312            r#"
313[listen."user-default"]
314enable = false
315[listen."system-default"]
316dir = "/usr/local/etc/arti-rpc/connect.d"
317file_options = { "tmp.toml" = { "enable" = false } }
318[listen."my-connpt"]
319file = "/home/dante/.paradiso/connpt.toml"
320"#,
321        )
322        .unwrap();
323        let mut altered_user_defaults = user_defaults_builder.clone();
324        altered_user_defaults.listener_options().enable(false);
325        let mut altered_system_defaults = system_defaults_builder.clone();
326        altered_system_defaults.dir(CfgPath::new(
327            "/usr/local/etc/arti-rpc/connect.d".to_string(),
328        ));
329        let mut opt = ConnectPointOptionsBuilder::default();
330        opt.enable(false);
331        altered_system_defaults
332            .file_options()
333            .insert("tmp.toml".to_string(), opt);
334        let mut my_connpt = RpcListenerSetConfigBuilder::default();
335        my_connpt.file(CfgPath::new(
336            "/home/dante/.paradiso/connpt.toml".to_string(),
337        ));
338
339        assert_eq!(
340            altered,
341            RpcConfig {
342                enable: false,
343                listen: vec![
344                    (
345                        "user-default".to_string(),
346                        altered_user_defaults.build().unwrap()
347                    ),
348                    (
349                        "system-default".to_string(),
350                        altered_system_defaults.build().unwrap()
351                    ),
352                    ("my-connpt".to_string(), my_connpt.build().unwrap()),
353                ]
354                .into_iter()
355                .collect(),
356                listen_default: listen_defaults_defaults()
357            }
358        );
359    }
360}