1
//! Experimental RPC support.
2

            
3
use anyhow::Result;
4
use arti_rpcserver::RpcMgr;
5
use derive_deftly::Deftly;
6
use fs_mistrust::Mistrust;
7
use futures::{AsyncReadExt, stream::StreamExt};
8
use session::ArtiRpcSession;
9
use std::collections::BTreeMap;
10
use std::{io::Result as IoResult, sync::Arc};
11
use tor_config::derive::prelude::*;
12
use tor_config_path::CfgPathResolver;
13
use tracing::{debug, info};
14

            
15
use arti_client::TorClient;
16
use tor_rtcompat::{NetStreamListener as _, Runtime, SpawnExt, general};
17

            
18
pub(crate) mod conntarget;
19
pub(crate) mod listener;
20
mod proxyinfo;
21
mod session;
22
mod superuser;
23

            
24
use listener::RpcListenerSetConfig;
25
pub(crate) use session::{RpcStateSender, RpcVisibleArtiState};
26

            
27
use crate::reload_cfg::{CfgMgr, LaunchableTorClient};
28
use crate::rpc::superuser::RpcSuperuser;
29

            
30
/// Configuration for Arti's RPC subsystem.
31
///
32
/// You cannot change this section on a running Arti client.
33
#[derive(Debug, Clone, Deftly, Eq, PartialEq)]
34
#[derive_deftly(TorConfig)]
35
#[cfg_attr(feature = "experimental-api", visibility::make(pub))]
36
#[cfg_attr(feature = "experimental-api", deftly(tor_config(vis = "pub")))]
37
pub(crate) struct RpcConfig {
38
    /// If true, then the RPC subsystem is enabled and will listen for connections.
39
    #[deftly(tor_config(default = "false"))] // TODO RPC make this true once we are stable.
40
    enable: bool,
41

            
42
    /// A set of named locations in which to find connect files.
43
    #[deftly(tor_config(map, default = "listener::listener_map_defaults()"))]
44
    listen: BTreeMap<String, RpcListenerSetConfig>,
45

            
46
    /// A list of default connect points to bind
47
    /// if no enabled connect points are found under `listen`.
48
    #[deftly(tor_config(list(element(clone)), default = "listen_defaults_defaults()"))]
49
    listen_default: Vec<String>,
50
}
51

            
52
/// Return default values for `RpcConfig.listen_default`
53
372
fn listen_defaults_defaults() -> Vec<String> {
54
372
    vec![tor_rpc_connect::USER_DEFAULT_CONNECT_POINT.to_string()]
55
372
}
56

            
57
/// Information about an incoming connection.
58
///
59
/// Yielded in a stream from our RPC listeners.
60
type IncomingConn = (
61
    general::Stream,
62
    general::SocketAddr,
63
    Arc<listener::RpcConnInfo>,
64
);
65

            
66
/// Bind to all configured RPC listeners in `cfg`.
67
///
68
/// On success, return a stream of `IncomingConn`.
69
async fn launch_all_listeners<R: Runtime>(
70
    runtime: &R,
71
    cfg: &RpcConfig,
72
    resolver: &CfgPathResolver,
73
    mistrust: &Mistrust,
74
) -> anyhow::Result<(
75
    impl futures::Stream<Item = IoResult<IncomingConn>> + Unpin + use<R>,
76
    Vec<tor_rpc_connect::server::ListenerGuard>,
77
)> {
78
    let mut listeners = Vec::new();
79
    let mut guards = Vec::new();
80
    for (name, listener_cfg) in cfg.listen.iter() {
81
        for (lis, info, guard) in listener_cfg
82
            .bind(runtime, name.as_str(), resolver, mistrust)
83
            .await?
84
        {
85
            // (Note that `bind` only returns enabled listeners, so we don't need to check here.
86
            debug!(
87
                "Listening at {} for {}",
88
                lis.local_addr()
89
                    .expect("general::listener without address?")
90
                    .display_lossy(),
91
                info.name,
92
            );
93
            listeners.push((lis, info));
94
            guards.push(guard);
95
        }
96
    }
97
    if listeners.is_empty() {
98
        for (idx, connpt) in cfg.listen_default.iter().enumerate() {
99
            let display_index = idx + 1; // One-indexed values are more human-readable.
100
            let (lis, info, guard) =
101
                listener::bind_string(connpt, display_index, runtime, resolver, mistrust).await?;
102
            debug!(
103
                "Listening at {} for {}",
104
                lis.local_addr()
105
                    .expect("general::listener without address?")
106
                    .display_lossy(),
107
                info.name,
108
            );
109
            listeners.push((lis, info));
110
            guards.push(guard);
111
        }
112
    }
113
    if listeners.is_empty() {
114
        info!("No RPC listeners configured.");
115
    }
116

            
117
    let streams = listeners.into_iter().map(|(listener, info)| {
118
        listener
119
            .incoming()
120
            .map(move |accept_result| match accept_result {
121
                Ok((netstream, addr)) => Ok((netstream, addr, Arc::clone(&info))),
122
                Err(e) => Err(e),
123
            })
124
    });
125

            
126
    Ok((futures::stream::select_all(streams), guards))
127
}
128

            
129
/// Create an RPC manager, bind to connect points, and open a listener task to accept incoming
130
/// RPC connections.
131
pub(crate) async fn launch_rpc_mgr<R: Runtime>(
132
    runtime: &R,
133
    cfg: &RpcConfig,
134
    resolver: &CfgPathResolver,
135
    mistrust: &Mistrust,
136
    client: Arc<TorClient<R>>,
137
    launchable: Arc<LaunchableTorClient<R>>,
138
    cfg_mgr: Arc<CfgMgr<R>>,
139
) -> Result<Option<RpcProxySupport>> {
140
    if !cfg.enable {
141
        return Ok(None);
142
    }
143
    let (rpc_state, rpc_state_sender) = RpcVisibleArtiState::new();
144

            
145
    let rpc_mgr = RpcMgr::new()?;
146
    // Register methods. Needed since TorClient is generic.
147
    //
148
    // TODO: If we accumulate a large number of generics like this, we should do this elsewhere.
149
    rpc_mgr.register_rpc_methods(TorClient::<R>::rpc_methods());
150
    rpc_mgr.register_rpc_methods(arti_rpcserver::rpc_methods::<R>());
151
    rpc_mgr.register_rpc_methods(RpcSuperuser::<R>::rpc_methods());
152

            
153
    let rt_clone = runtime.clone();
154
    let rpc_mgr_clone = rpc_mgr.clone();
155

            
156
    let (incoming, guards) = launch_all_listeners(runtime, cfg, resolver, mistrust).await?;
157

            
158
    // TODO: Using spawn in this way makes it hard to report whether we
159
    // succeeded or not. This is something we should fix when we refactor
160
    // our service-launching code.
161
    runtime.spawn(async move {
162
        let result = run_rpc_listener(
163
            rt_clone,
164
            incoming,
165
            rpc_mgr_clone,
166
            client,
167
            launchable,
168
            cfg_mgr,
169
            rpc_state,
170
        )
171
        .await;
172
        if let Err(e) = result {
173
            tracing::warn!("RPC manager quit with an error: {}", e);
174
        }
175
        drop(guards);
176
    })?;
177
    Ok(Some(RpcProxySupport {
178
        rpc_mgr,
179
        rpc_state_sender,
180
    }))
181
}
182

            
183
/// Backend function to implement an RPC listener: runs in a loop.
184
async fn run_rpc_listener<R: Runtime>(
185
    runtime: R,
186
    mut incoming: impl futures::Stream<Item = IoResult<IncomingConn>> + Unpin,
187
    rpc_mgr: Arc<RpcMgr>,
188
    client: Arc<TorClient<R>>,
189
    launchable: Arc<LaunchableTorClient<R>>,
190
    cfg_mgr: Arc<CfgMgr<R>>,
191
    rpc_state: Arc<RpcVisibleArtiState>,
192
) -> Result<()> {
193
    while let Some((stream, _addr, info)) = incoming.next().await.transpose()? {
194
        debug!("Received incoming RPC connection from {}", &info.name);
195

            
196
        let client_clone = client.clone();
197
        let rpc_state_clone = rpc_state.clone();
198
        let launchable = launchable.clone();
199
        let cfg_mgr_clone = cfg_mgr.clone();
200
        let connection = rpc_mgr.new_connection(info.auth.clone(), move |auth| {
201
            ArtiRpcSession::new(
202
                auth,
203
                &client_clone,
204
                &launchable,
205
                &rpc_state_clone,
206
                &cfg_mgr_clone,
207
                &info,
208
            ) as _
209
        });
210
        let (input, output) = stream.split();
211

            
212
        runtime.spawn(async {
213
            let result = connection.run(input, output).await;
214
            if let Err(e) = result {
215
                tracing::warn!("RPC session ended with an error: {}", e);
216
            }
217
        })?;
218
    }
219
    Ok(())
220
}
221

            
222
/// Information passed to a proxy or similar stream provider when running with RPC support.
223
pub(crate) struct RpcProxySupport {
224
    /// An RPC manager to use for looking up objects as possible stream targets.
225
    pub(crate) rpc_mgr: Arc<arti_rpcserver::RpcMgr>,
226
    /// An RPCStateSender to use for registering the list of known proxy ports.
227
    pub(crate) rpc_state_sender: RpcStateSender,
228
}
229

            
230
#[cfg(test)]
231
mod test {
232
    // @@ begin test lint list maintained by maint/add_warning @@
233
    #![allow(clippy::bool_assert_comparison)]
234
    #![allow(clippy::clone_on_copy)]
235
    #![allow(clippy::dbg_macro)]
236
    #![allow(clippy::mixed_attributes_style)]
237
    #![allow(clippy::print_stderr)]
238
    #![allow(clippy::print_stdout)]
239
    #![allow(clippy::single_char_pattern)]
240
    #![allow(clippy::unwrap_used)]
241
    #![allow(clippy::unchecked_time_subtraction)]
242
    #![allow(clippy::useless_vec)]
243
    #![allow(clippy::needless_pass_by_value)]
244
    #![allow(clippy::string_slice)] // See arti#2571
245
    //! <!-- @@ end test lint list maintained by maint/add_warning @@ -->
246

            
247
    use listener::{ConnectPointOptionsBuilder, RpcListenerSetConfigBuilder};
248
    use tor_config_path::CfgPath;
249
    use tor_rpc_connect::ParsedConnectPoint;
250

            
251
    use super::*;
252

            
253
    #[test]
254
    fn rpc_method_names() {
255
        // We run this from a nice high level module, to ensure that as many method names as
256
        // possible will be in-scope.
257
        let problems = tor_rpcbase::check_method_names([]);
258

            
259
        for (m, err) in &problems {
260
            eprintln!("Bad method name {m:?}: {err}");
261
        }
262
        assert!(problems.is_empty());
263
    }
264

            
265
    #[test]
266
    fn parse_listener_defaults() {
267
        for string in listen_defaults_defaults() {
268
            let _parsed: ParsedConnectPoint = string.parse().unwrap();
269
        }
270
    }
271

            
272
    #[test]
273
    fn parsing_and_building() {
274
        fn build(s: &str) -> Result<RpcConfig, anyhow::Error> {
275
            let b: RpcConfigBuilder = toml::from_str(s)?;
276
            Ok(b.build()?)
277
        }
278

            
279
        let mut user_defaults_builder = RpcListenerSetConfigBuilder::default();
280
        user_defaults_builder.listener_options().enable(true);
281
        user_defaults_builder.dir(CfgPath::new("${ARTI_LOCAL_DATA}/rpc/connect.d".to_string()));
282
        let mut system_defaults_builder = RpcListenerSetConfigBuilder::default();
283
        system_defaults_builder.listener_options().enable(false);
284
        system_defaults_builder.dir(CfgPath::new("/etc/arti-rpc/connect.d".to_string()));
285

            
286
        // Make sure that an empty configuration gets us the defaults.
287
        let defaults = build("").unwrap();
288
        assert_eq!(
289
            defaults,
290
            RpcConfig {
291
                enable: false,
292
                listen: vec![
293
                    (
294
                        "user-default".to_string(),
295
                        user_defaults_builder.build().unwrap()
296
                    ),
297
                    (
298
                        "system-default".to_string(),
299
                        system_defaults_builder.build().unwrap()
300
                    ),
301
                ]
302
                .into_iter()
303
                .collect(),
304
                listen_default: listen_defaults_defaults()
305
            }
306
        );
307

            
308
        // Make sure that overriding specific options works as expected.
309
        let altered = build(
310
            r#"
311
[listen."user-default"]
312
enable = false
313
[listen."system-default"]
314
dir = "/usr/local/etc/arti-rpc/connect.d"
315
file_options = { "tmp.toml" = { "enable" = false } }
316
[listen."my-connpt"]
317
file = "/home/dante/.paradiso/connpt.toml"
318
"#,
319
        )
320
        .unwrap();
321
        let mut altered_user_defaults = user_defaults_builder.clone();
322
        altered_user_defaults.listener_options().enable(false);
323
        let mut altered_system_defaults = system_defaults_builder.clone();
324
        altered_system_defaults.dir(CfgPath::new(
325
            "/usr/local/etc/arti-rpc/connect.d".to_string(),
326
        ));
327
        let mut opt = ConnectPointOptionsBuilder::default();
328
        opt.enable(false);
329
        altered_system_defaults
330
            .file_options()
331
            .insert("tmp.toml".to_string(), opt);
332
        let mut my_connpt = RpcListenerSetConfigBuilder::default();
333
        my_connpt.file(CfgPath::new(
334
            "/home/dante/.paradiso/connpt.toml".to_string(),
335
        ));
336

            
337
        assert_eq!(
338
            altered,
339
            RpcConfig {
340
                enable: false,
341
                listen: vec![
342
                    (
343
                        "user-default".to_string(),
344
                        altered_user_defaults.build().unwrap()
345
                    ),
346
                    (
347
                        "system-default".to_string(),
348
                        altered_system_defaults.build().unwrap()
349
                    ),
350
                    ("my-connpt".to_string(), my_connpt.build().unwrap()),
351
                ]
352
                .into_iter()
353
                .collect(),
354
                listen_default: listen_defaults_defaults()
355
            }
356
        );
357
    }
358
}