1use 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#[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 #[deftly(tor_config(default = false))] enable: bool,
43
44 #[deftly(tor_config(map, default = listener::listener_map_defaults()))]
46 listen: BTreeMap<String, RpcListenerSetConfig>,
47
48 #[deftly(tor_config(list(element(clone)), default = listen_defaults_defaults()))]
51 listen_default: Vec<String>,
52}
53
54fn listen_defaults_defaults() -> Vec<String> {
56 vec![tor_rpc_connect::USER_DEFAULT_CONNECT_POINT.to_string()]
57}
58
59type IncomingConn = (
63 general::Stream,
64 general::SocketAddr,
65 Arc<listener::RpcConnInfo>,
66);
67
68async 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 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; 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
131pub(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 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 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
185async 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
224pub(crate) struct RpcProxySupport {
226 pub(crate) rpc_mgr: Arc<arti_rpcserver::RpcMgr>,
228 pub(crate) rpc_state_sender: RpcStateSender,
230}
231
232#[cfg(test)]
233mod test {
234 #![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)] 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 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 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 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}