1#![cfg_attr(docsrs, feature(doc_cfg))]
2#![doc = include_str!("../README.md")]
3#![allow(renamed_and_removed_lints)] #![allow(unknown_lints)] #![warn(missing_docs)]
7#![warn(noop_method_call)]
8#![warn(unreachable_pub)]
9#![warn(clippy::all)]
10#![deny(clippy::await_holding_lock)]
11#![deny(clippy::cargo_common_metadata)]
12#![deny(clippy::cast_lossless)]
13#![deny(clippy::checked_conversions)]
14#![allow(clippy::cognitive_complexity)] #![deny(clippy::debug_assert_with_mut_call)]
16#![deny(clippy::exhaustive_enums)]
17#![deny(clippy::exhaustive_structs)]
18#![deny(clippy::expl_impl_clone_on_copy)]
19#![deny(clippy::fallible_impl_from)]
20#![deny(clippy::implicit_clone)]
21#![deny(clippy::large_stack_arrays)]
22#![warn(clippy::manual_ok_or)]
23#![deny(clippy::missing_docs_in_private_items)]
24#![warn(clippy::needless_borrow)]
25#![warn(clippy::needless_pass_by_value)]
26#![warn(clippy::option_option)]
27#![deny(clippy::print_stderr)]
28#![deny(clippy::print_stdout)]
29#![warn(clippy::rc_buffer)]
30#![deny(clippy::ref_option_ref)]
31#![warn(clippy::semicolon_if_nothing_returned)]
32#![warn(clippy::trait_duplication_in_bounds)]
33#![deny(clippy::unchecked_time_subtraction)]
34#![deny(clippy::unnecessary_wraps)]
35#![warn(clippy::unseparated_literal_suffix)]
36#![deny(clippy::unwrap_used)]
37#![deny(clippy::mod_module_files)]
38#![allow(clippy::let_unit_value)] #![allow(clippy::uninlined_format_args)]
40#![allow(clippy::significant_drop_in_scrutinee)] #![allow(clippy::result_large_err)] #![allow(clippy::needless_raw_string_hashes)] #![allow(clippy::needless_lifetimes)] #![allow(mismatched_lifetime_syntaxes)] #![allow(clippy::collapsible_if)] #![deny(clippy::unused_async)]
47#![deny(clippy::string_slice)] pub mod config;
51pub mod err;
52
53#[cfg(feature = "managed-pts")]
54pub mod ipc;
55
56#[cfg(feature = "managed-pts")]
57mod managed;
58
59use crate::config::{TransportConfig, TransportOptions};
60use crate::err::PtError;
61use std::collections::HashMap;
62use std::net::SocketAddr;
63use std::path::PathBuf;
64use std::sync::{Arc, RwLock};
65use tor_chanmgr::ProxyProtocol;
66use tor_config_path::CfgPathResolver;
67use tor_linkspec::PtTransportName;
68use tor_rtcompat::Runtime;
69use tor_socksproto::SocksVersion;
70use tracing::warn;
71#[cfg(feature = "managed-pts")]
72use {
73 crate::managed::{PtReactor, PtReactorMessage},
74 futures::channel::mpsc::{self, UnboundedSender},
75 tor_error::error_report,
76 tor_rtcompat::SpawnExt,
77};
78#[cfg(feature = "tor-channel-factory")]
79use {
80 async_trait::async_trait,
81 tor_chanmgr::{
82 builder::ChanBuilder,
83 factory::{AbstractPtError, ChannelFactory},
84 transport::ExternalProxyPlugin,
85 },
86 tracing::trace,
87};
88#[cfg(all(feature = "managed-pts", feature = "tor-channel-factory"))]
89use {oneshot_fused_workaround as oneshot, tracing::info};
90
91#[derive(Default, Debug)]
93struct PtSharedState {
94 #[cfg(feature = "managed-pts")]
98 managed_cmethods: HashMap<PtTransportName, PtClientMethod>,
99 configured: HashMap<PtTransportName, TransportOptions>,
101 outbound_proxy: Option<ProxyProtocol>,
103}
104
105pub struct PtMgr<R> {
108 #[allow(dead_code)]
110 runtime: R,
111 state: Arc<RwLock<PtSharedState>>,
113 #[cfg(feature = "managed-pts")]
115 tx: UnboundedSender<PtReactorMessage>,
116}
117
118impl<R: Runtime> PtMgr<R> {
119 fn transform_config(
121 binaries: Vec<TransportConfig>,
122 ) -> Result<HashMap<PtTransportName, TransportOptions>, tor_error::Bug> {
123 let mut ret = HashMap::new();
124 for thing in binaries {
129 for tn in thing.protocols.iter() {
130 ret.insert(tn.clone(), thing.clone().try_into()?);
131 }
132 }
133 for opt in ret.values() {
134 match opt {
135 TransportOptions::Unmanaged(u) => {
136 if !u.is_localhost() {
137 warn!(
138 "Configured to connect to a PT on a non-local addresses. This is usually insecure! We recommend running PTs on localhost only."
139 );
140 }
141 }
142 #[cfg(feature = "managed-pts")]
143 TransportOptions::Managed(_) => {
144 }
148 }
149 }
150 Ok(ret)
151 }
152
153 pub fn new(
156 transports: Vec<TransportConfig>,
157 #[allow(unused)] state_dir: PathBuf,
158 #[allow(unused)] path_resolver: Arc<CfgPathResolver>,
159 outbound_proxy: Option<ProxyProtocol>,
160 rt: R,
161 ) -> Result<Self, PtError> {
162 let state = PtSharedState {
163 #[cfg(feature = "managed-pts")]
164 managed_cmethods: Default::default(),
165 configured: Self::transform_config(transports)?,
166 outbound_proxy,
167 };
168 let state = Arc::new(RwLock::new(state));
169
170 #[cfg(feature = "managed-pts")]
172 let tx = {
173 let (tx, rx) = mpsc::unbounded();
174
175 let mut reactor =
176 PtReactor::new(rt.clone(), state.clone(), rx, state_dir, path_resolver);
177 rt.spawn(async move {
178 loop {
179 match reactor.run_one_step().await {
180 Ok(true) => return,
181 Ok(false) => {}
182 Err(e) => {
183 error_report!(e, "PtReactor failed");
184 return;
185 }
186 }
187 }
188 })
189 .map_err(|e| PtError::Spawn { cause: Arc::new(e) })?;
190
191 tx
192 };
193
194 Ok(Self {
195 runtime: rt,
196 state,
197 #[cfg(feature = "managed-pts")]
198 tx,
199 })
200 }
201
202 pub fn reconfigure(
204 &self,
205 how: tor_config::Reconfigure,
206 transports: Vec<TransportConfig>,
207 outbound_proxy: Option<ProxyProtocol>,
208 ) -> Result<(), tor_config::ReconfigureError> {
209 let configured = Self::transform_config(transports)?;
210 if how == tor_config::Reconfigure::CheckAllOrNothing {
211 return Ok(());
212 }
213 {
214 let mut inner = self.state.write().expect("ptmgr poisoned");
215 inner.configured = configured;
216 inner.outbound_proxy = outbound_proxy;
217 }
218 #[cfg(feature = "managed-pts")]
225 let _ = self.tx.unbounded_send(PtReactorMessage::Reconfigured);
226 Ok(())
227 }
228
229 #[cfg(feature = "tor-channel-factory")]
235 async fn get_cmethod_for_transport(
236 &self,
237 transport: &PtTransportName,
238 ) -> Result<Option<PtClientMethod>, PtError> {
239 let (cfg, managed_cmethod) = {
240 let inner = self.state.read().expect("ptmgr poisoned");
244 let cfg = inner.configured.get(transport);
245 let managed_cmethod = inner.managed_cmethods.get(transport);
246 (cfg.cloned(), managed_cmethod.cloned())
247 };
248
249 #[cfg(not(feature = "managed-pts"))]
250 let _ = managed_cmethod; match cfg {
253 Some(TransportOptions::Unmanaged(cfg)) => {
254 let cmethod = cfg.cmethod();
255 trace!(
256 "Found configured unmanaged transport {transport} accessible via {cmethod:?}"
257 );
258 Ok(Some(cmethod))
259 }
260 #[cfg(feature = "managed-pts")]
261 Some(TransportOptions::Managed(_cfg)) => {
262 match managed_cmethod {
263 Some(cmethod) => {
265 trace!(
266 "Found configured managed transport {transport} accessible via {cmethod:?}"
267 );
268 Ok(Some(cmethod))
269 }
270 None => {
272 Ok(Some(self.spawn_transport(transport).await?))
294 }
295 }
296 }
297 None => {
299 trace!("Got a request for transport {transport}, which is not configured.");
300 Ok(None)
301 }
302 }
303 }
304
305 #[cfg(all(feature = "tor-channel-factory", feature = "managed-pts"))]
307 async fn spawn_transport(
308 &self,
309 transport: &PtTransportName,
310 ) -> Result<PtClientMethod, PtError> {
311 info!(
314 "Got a request for transport {transport}, which is not currently running. Launching it."
315 );
316
317 let (tx, rx) = oneshot::channel();
318 self.tx
319 .unbounded_send(PtReactorMessage::Spawn {
320 pt: transport.clone(),
321 result: tx,
322 })
323 .map_err(|_| {
324 PtError::Internal(tor_error::internal!("PT reactor closed unexpectedly"))
325 })?;
326
327 let method = match rx.await {
328 Err(_) => {
329 return Err(PtError::Internal(tor_error::internal!(
330 "PT reactor closed unexpectedly"
331 )));
332 }
333 Ok(Err(e)) => {
334 warn!("PT for {transport} failed to launch: {e}");
335 return Err(e);
336 }
337 Ok(Ok(method)) => method,
338 };
339
340 info!("Successfully launched PT for {transport} at {method:?}.");
341 Ok(method)
342 }
343}
344
345#[derive(Debug, Clone, PartialEq, Eq)]
347pub struct PtClientMethod {
348 pub(crate) kind: SocksVersion,
350 pub(crate) endpoint: SocketAddr,
352}
353
354impl PtClientMethod {
355 pub fn kind(&self) -> SocksVersion {
357 self.kind
358 }
359
360 pub fn endpoint(&self) -> SocketAddr {
362 self.endpoint
363 }
364}
365
366#[cfg(feature = "tor-channel-factory")]
367#[async_trait]
368impl<R: Runtime> tor_chanmgr::factory::AbstractPtMgr for PtMgr<R> {
369 async fn factory_for_transport(
370 &self,
371 transport: &PtTransportName,
372 ) -> Result<Option<Arc<dyn ChannelFactory + Send + Sync>>, Arc<dyn AbstractPtError>> {
373 let cmethod = match self.get_cmethod_for_transport(transport).await {
374 Err(e) => return Err(Arc::new(e)),
375 Ok(None) => return Ok(None),
376 Ok(Some(m)) => m,
377 };
378
379 let proxy = ExternalProxyPlugin::new(self.runtime.clone(), cmethod.endpoint, cmethod.kind);
380 let factory = ChanBuilder::new_client(self.runtime.clone(), proxy);
381 Ok(Some(Arc::new(factory)))
384 }
385}