Skip to main content

arti/
reload_cfg.rs

1//! Code to watch configuration files for any changes.
2
3use std::collections::HashSet;
4use std::sync::{Arc, Mutex, Weak};
5use std::time::Duration;
6
7use anyhow::Context;
8use arti_client::TorClient;
9use arti_client::config::Reconfigure;
10use futures::StreamExt;
11use futures::stream::BoxStream;
12use futures::{FutureExt as _, Stream, select_biased};
13use tor_basic_utils::error_sources::ErrorSources;
14use tor_config::ReconfigureError;
15use tor_config::file_watcher::{
16    self, FileEventReceiver, FileEventSender, FileWatcher, FileWatcherBuilder,
17};
18use tor_config::load::{ConfigResolveOptions, DisfavouredKey};
19use tor_config::{ConfigGetValueError, ConfigurationTree};
20use tor_config::{ConfigurationSource, ConfigurationSources, sources::FoundConfigFiles};
21use tor_error::warn_report;
22use tor_error::{HasKind, into_internal};
23use tor_rtcompat::Runtime;
24use tor_rtcompat::SpawnExt;
25use tracing::{debug, error, info, instrument, warn};
26
27#[cfg(target_family = "unix")]
28use crate::process::sighup_stream;
29#[cfg(feature = "rpc")]
30use crate::rpc;
31
32#[cfg(not(target_family = "unix"))]
33use futures::stream;
34
35use crate::{ArtiCombinedConfig, ArtiConfig};
36
37/// How long to wait after an event got received, before we try to process it.
38const DEBOUNCE_INTERVAL: Duration = Duration::from_secs(1);
39
40/// An object that can be reconfigured when our configuration changes.
41///
42/// We use this trait so that we can represent abstract modules in our
43/// application, and pass the configuration to each of them.
44//
45// TODO: It is very likely we will want to refactor this even further once we
46// have a notion of what our modules truly are.
47#[cfg_attr(feature = "experimental-api", visibility::make(pub))]
48pub(crate) trait ReconfigurableModule: Send + Sync {
49    /// Try to reconfigure this module according to a newly loaded configuration.
50    ///
51    /// See [`Reconfigure`] for a description of error-handling behavior.
52    fn reconfigure(
53        &self,
54        new: &ArtiCombinedConfig,
55        how: Reconfigure,
56    ) -> Result<(), ReconfigureError>;
57}
58
59/// Structure to reload configuration as necessary.
60#[cfg_attr(feature = "experimental-api", visibility::make(pub))]
61pub(crate) struct CfgMgr<R> {
62    /// A runtime that we use when constructing [`FileWatcher`]s.
63    runtime: R,
64
65    /// The sources from which we read our configuration.
66    sources: ConfigurationSources,
67
68    /// A sender to use when constructing new [`FileWatcher`]s.
69    tx: FileEventSender,
70
71    /// Mutable state.
72    inner: Mutex<CfgMgrInner>,
73}
74
75/// Mutable part of a CfgMgr.
76///
77/// ## RPC Dataflow
78///
79/// We keep a fair amount of state for our configuration,
80/// especially when we are supporting RPC.  A few important pieces are,
81/// in dataflow order:
82///
83/// 1. `sources`: A set of [`ConfigurationSources`] telling us where to load
84///    our configuration from.
85/// 2. `loaded_cfg`: A [`ConfigurationTree`] that we have loaded from our
86///    `sources`.  This is the most recent tree that we were able to successfully
87///    decode and apply. (RPC only)
88/// 3. `additional_cfg`: A set of [`rpc::ConfigSettings`] provided by an RPC superuser app,
89///    to override options in `loaded_cfg`. (RPC only)
90/// 4. `normalized_cfg`: A [`ConfigurationTree`] produced by combining
91///    loaded_cfg` and `additional_cfg`, and filling any missing defaults.
92///
93/// When we are reloading our configuration from disk, we use `sources` to
94/// load and parse a new [`ConfigurationTree`], then we apply `additional_cfg` to it,
95/// and then we see whether that configuration can successfully be resolved
96/// and used to reconfigure the modules in Arti.
97/// On success, we replace `loaded_cfg` and `normalized_cfg`,
98///
99/// When we are changing our configuration via RPC, we try to change `additional_cfg`,
100/// apply it to `loaded_cfg`,
101/// and then we see whether that configuration can successfully be resolved
102/// and used to reconfigure the modules in Arti.
103/// On success, we replace `normalized_cfg` and `additional_cfg`.
104#[derive(Default)]
105struct CfgMgrInner {
106    /// A list of modules to alert whenever the configuration has changed.
107    modules: Vec<Weak<dyn ReconfigurableModule>>,
108
109    /// If present, a [`FileWatcher`] that is currently watching for changes
110    /// in the configuration files and directories.
111    watcher: Option<FileWatcher>,
112
113    /// RPC only: The most recent configuration tree _as loaded_.  We use this to apply
114    /// additional_cfg repeatedly without reloading all the files every time RPC tells us
115    /// to change something.
116    #[cfg(feature = "rpc")]
117    loaded_cfg: ConfigurationTree,
118
119    /// RPC only: A set of RPC-provided options to apply to the configuration before decoding it.
120    #[cfg(feature = "rpc")]
121    additional_cfg: rpc::ConfigSettings,
122
123    /// RPC only: a fully populated, normalized configuration tree, based on the most recent time
124    /// that we called [`CfgMgr::reload_configuration`].
125    #[cfg(feature = "rpc")]
126    normalized_cfg: ConfigurationTree,
127
128    /// RPC only: a set of unrecognized options from the configuration.
129    #[cfg(feature = "rpc")]
130    unrecognized_keys: HashSet<DisfavouredKey>,
131
132    /// RPC only: a set of deprecated options from the configuration
133    #[cfg(feature = "rpc")]
134    deprecated_keys: HashSet<DisfavouredKey>,
135}
136
137/// A watcher process that we have not yet launched.
138#[cfg_attr(feature = "experimental-api", visibility::make(pub))]
139#[must_use = "UnlaunchedWatcher does nothing unless you launch it."]
140pub(crate) struct UnlaunchedWatcher<R> {
141    /// The related [`CfgMgr`] that we should tell about reconfiguration events.
142    weak_mgr: Weak<CfgMgr<R>>,
143
144    /// A stream on which we will get alerts about SIGHUP events.
145    sighup_stream: BoxStream<'static, ()>,
146
147    /// A stream that will tell us when our files are changed.
148    watcher_rx: FileEventReceiver,
149
150    /// An interval that we wait to debounce events from watcher_rx or sighup_stream.
151    debounce_interval: Option<Duration>,
152
153    /// If true, we start watching for file changes immediately at launch.
154    watch_files_at_start: bool,
155}
156
157impl<R: Runtime> CfgMgr<R> {
158    /// Construct a new CfgMgr, and launch a task to watch for any events
159    /// that mean we have to reload our configuration.
160    ///
161    /// If the provided configuration requires it, watch for changes in `sources`
162    /// and try to reload our configuration. On unix platforms, also watch
163    /// for SIGHUP and reload configuration then.
164    ///
165    /// The modules are `Weak` references to prevent this background task
166    /// from keeping them alive.
167    ///
168    /// See the [`FileWatcher`](FileWatcher#Limitations) docs for limitations.
169    #[cfg_attr(feature = "experimental-api", visibility::make(pub))]
170    #[instrument(level = "trace", skip_all)]
171    pub(crate) fn new(
172        runtime: R,
173        sources: ConfigurationSources,
174        #[cfg(feature = "rpc")] loaded_cfg: ConfigurationTree,
175        config: &ArtiConfig,
176        modules: Vec<Weak<dyn ReconfigurableModule>>,
177    ) -> anyhow::Result<(Arc<Self>, UnlaunchedWatcher<R>)> {
178        let (tx, rx) = file_watcher::channel();
179        let mgr = Arc::new(CfgMgr {
180            runtime,
181            sources,
182            tx,
183            inner: Mutex::new(CfgMgrInner {
184                modules,
185                #[cfg(feature = "rpc")]
186                loaded_cfg,
187                ..Default::default()
188            }),
189        });
190
191        cfg_if::cfg_if! {
192            if #[cfg(target_family = "unix")] {
193                let sighup_stream = sighup_stream()?;
194            } else {
195                let sighup_stream = stream::pending();
196            }
197        }
198        let sighup_stream = sighup_stream.boxed();
199
200        let watcher = UnlaunchedWatcher {
201            weak_mgr: Arc::downgrade(&mgr),
202            sighup_stream,
203            watcher_rx: rx,
204            debounce_interval: Some(DEBOUNCE_INTERVAL),
205            watch_files_at_start: config.application().watch_configuration,
206        };
207
208        Ok((mgr, watcher))
209    }
210
211    /// Create a new [`FileWatcher`] for the files in this configuration.
212    ///
213    /// Return it, along with the set of files we found.
214    ///
215    /// The caller is responsible for storing the `FileWatcher`; when it is dropped,
216    /// it stops watching.
217    fn launch_file_watcher(&self) -> anyhow::Result<(FileWatcher, FoundConfigFiles<'_>)> {
218        let mut watcher = FileWatcher::builder(self.runtime.clone());
219        let found_files = prepare(&mut watcher, &self.sources)?;
220        let watcher = watcher.start_watching(self.tx.clone())?;
221        Ok((watcher, found_files))
222    }
223
224    /// Reload the configuration.
225    #[instrument(level = "trace", skip_all)]
226    #[cfg_attr(feature = "experimental-api", visibility::make(pub))]
227    pub(crate) fn reload_configuration(&self, how: Reconfigure) -> anyhow::Result<()> {
228        let mut inner = self.inner.lock().expect("Lock poisoned");
229
230        // Question: I do not understand why we are making a new file watcher unconditionally
231        // at this point. -nm
232        let (found_files, new_watcher) = if inner.watcher.is_some() {
233            let (watcher, files) = self
234                .launch_file_watcher()
235                .context("Failed to re-scan config")?;
236            (files, Some(watcher))
237        } else {
238            let files = self
239                .sources
240                .scan()
241                .context("FS watch: failed to rescan config")?;
242            (files, None)
243        };
244
245        let mut config = found_files.load()?;
246        #[cfg(feature = "rpc")]
247        let loaded_cfg = config.clone();
248        #[cfg(feature = "rpc")]
249        config.merge_from(&inner.additional_cfg)?;
250
251        match reconfigure(&config, &mut inner, how) {
252            Ok(watch) => {
253                info!("Successfully reloaded configuration.");
254                if how != Reconfigure::CheckAllOrNothing {
255                    #[cfg(feature = "rpc")]
256                    {
257                        inner.loaded_cfg = loaded_cfg;
258                    }
259                    if watch && inner.watcher.is_none() {
260                        info!("Starting watching over configuration.");
261                        let (watcher, _files) = self
262                            .launch_file_watcher()
263                            .context("Starting to watch over config")?;
264                        inner.watcher = Some(watcher);
265                    } else if !watch && inner.watcher.is_some() {
266                        info!("Stopped watching over configuration.");
267                        inner.watcher = None;
268                    } else {
269                        inner.watcher = new_watcher;
270                    }
271                }
272            }
273            Err(e) => warn_report!(e, "Couldn't reload configuration"),
274        }
275
276        Ok(())
277    }
278
279    /// Try to change the set of RPC configuration options.
280    ///
281    /// On success, the configuration is changed, and the changes are applied.
282    ///
283    /// On failure, the configuration is not changed, and the changes are not applied.
284    #[cfg(feature = "rpc")]
285    #[instrument(level = "trace", skip_all)]
286    #[cfg_attr(feature = "experimental-api", visibility::make(pub))]
287    pub(crate) fn try_modify_cfg<F>(
288        &self,
289        func: F,
290        how: Reconfigure,
291    ) -> Result<(), ChangeConfigurationError>
292    where
293        F: FnOnce(&mut rpc::ConfigSettings) -> Result<(), ChangeConfigurationError>,
294    {
295        let mut inner = self.inner.lock().expect("Lock poisoned");
296
297        let mut new_additional = inner.additional_cfg.clone();
298        func(&mut new_additional)?;
299        let mut new_cfg = inner.loaded_cfg.clone();
300        new_cfg.merge_from(&new_additional)?;
301
302        let watch = reconfigure(&new_cfg, &mut inner, how)?;
303
304        if how != Reconfigure::CheckAllOrNothing {
305            // If we reached here, we were successful. Remember new_additional...
306            inner.additional_cfg = new_additional;
307
308            // And adjust the file watcher.
309            if !watch && inner.watcher.is_some() {
310                inner.watcher = None;
311            } else if watch && inner.watcher.is_none() {
312                match self.launch_file_watcher() {
313                    Ok((watcher, _)) => inner.watcher = Some(watcher),
314                    Err(e) => warn_report!(e, "Unable to launch file watcher"),
315                }
316            }
317        }
318
319        Ok(())
320    }
321
322    /// Return the configuration value for a given key, if any is set.
323    #[cfg(feature = "rpc")]
324    pub(crate) fn get_cfg_setting(
325        &self,
326        key: &str,
327    ) -> Result<Option<rpc::ConfigValue>, ConfigGetValueError> {
328        let settings: Option<rpc::ConfigValue> = self
329            .inner
330            .lock()
331            .expect("Lock poisoned")
332            .normalized_cfg
333            .get_serde_value(key)?;
334
335        Ok(settings)
336    }
337}
338
339impl<R: Runtime> UnlaunchedWatcher<R> {
340    /// Begin running the file watcher task for a given configuration manager.
341    #[cfg_attr(feature = "experimental-api", visibility::make(pub))]
342    #[instrument(level = "trace", skip_all)]
343    pub(crate) fn launch(self) -> anyhow::Result<()> {
344        let UnlaunchedWatcher {
345            weak_mgr,
346            sighup_stream,
347            watcher_rx,
348            debounce_interval,
349            watch_files_at_start,
350        } = self;
351        let Some(mgr) = weak_mgr.upgrade() else {
352            return Err(anyhow::anyhow!(
353                "CfgMgr disappeared before we could launch the monitor task"
354            ));
355        };
356
357        let rt = mgr.runtime.clone();
358        let weak_mgr = Arc::downgrade(&mgr);
359        mgr.runtime
360            .spawn(async move {
361                let res: anyhow::Result<()> =
362                    run_watcher(rt, watcher_rx, sighup_stream, weak_mgr, debounce_interval).await;
363                match res {
364                    Ok(()) => debug!("Config watcher task exiting"),
365                    // TODO: warn_report does not work on anyhow::Error.
366                    Err(e) => error!("Config watcher task exiting: {}", tor_error::Report(e)),
367                }
368            })
369            .context("failed to spawn task")?;
370
371        if watch_files_at_start {
372            // Note: You might think that there was a race condition here, where launching the
373            // watcher _now_ would fail to catch any file changes that had happened between
374            // reading the configuration initially and now.
375            //
376            // You'd be right, except that the [`FileWatcher`] code starts every new FileWatcher
377            // with a pending `rescan` event.
378            let (watcher, _files) = mgr.launch_file_watcher()?;
379            mgr.inner.lock().expect("lock poisoned").watcher = Some(watcher);
380        }
381
382        Ok(())
383    }
384
385    /// Add `module` to the set of modules that need to be reconfigured when the configuration changes.
386    ///
387    /// This method is on the [`UnlaunchedWatcher`] because is not (yet) meant to be called after
388    /// the watcher task is launched.
389    #[cfg_attr(feature = "experimental-api", visibility::make(pub))]
390    pub(crate) fn add_module(&self, module: &Arc<dyn ReconfigurableModule>) -> anyhow::Result<()> {
391        let weak_module = Arc::downgrade(module);
392
393        let Some(mgr) = self.weak_mgr.upgrade() else {
394            return Err(anyhow::anyhow!(
395                "CfgMgr disappeared before launching watcher task."
396            ));
397        };
398
399        let mut inner = mgr.inner.lock().expect("poisoned lock");
400        inner.modules.push(weak_module);
401        Ok(())
402    }
403}
404
405/// Start watching for configuration changes.
406///
407/// Spawned from [`UnlaunchedWatcher::launch`].
408#[instrument(level = "trace", skip_all)]
409async fn run_watcher<R: Runtime>(
410    runtime: R,
411    mut rx: FileEventReceiver,
412    mut sighup_stream: impl Stream<Item = ()> + Unpin,
413    weak_mgr: Weak<CfgMgr<R>>,
414    debounce_interval: Option<Duration>,
415) -> anyhow::Result<()> {
416    debug!("Entering FS event loop");
417
418    loop {
419        select_biased! {
420            event = sighup_stream.next().fuse() => {
421                let Some(()) = event else {
422                    break;
423                };
424
425                info!("Received SIGHUP");
426            },
427            event = rx.next().fuse() => {
428                if let Some(debounce_interval) = debounce_interval {
429                    runtime.sleep(debounce_interval).await;
430                }
431
432                while let Some(_ignore) = rx.try_recv() {
433                    // Discard other events, so that we only reload once.
434                    //
435                    // We can afford to treat both error cases from try_recv [Empty
436                    // and Disconnected] as meaning that we've discarded other
437                    // events: if we're disconnected, we'll notice it when we next
438                    // call recv() in the outer loop.
439                }
440                debug!("Config reload event {:?}: reloading configuration.", event);
441            },
442        }
443
444        if let Some(mgr) = weak_mgr.upgrade() {
445            mgr.reload_configuration(Reconfigure::WarnOnFailures)?;
446            drop(mgr);
447        } else {
448            debug!("Configuration mgr disappeared; exiting loop");
449            break;
450        }
451    }
452
453    Ok(())
454}
455
456/// A TorClient that we may or may not have told to start bootstrapping.
457pub(crate) struct LaunchableTorClient<R: Runtime> {
458    /// Original value of defer_bootstrap.
459    orig_defer_bootstrap: bool,
460
461    /// True if we have launched bootstrapping on the the client.
462    have_launched: Mutex<bool>,
463
464    /// The client itself.
465    client: Arc<TorClient<R>>,
466}
467
468impl<R: Runtime> ReconfigurableModule for LaunchableTorClient<R> {
469    #[instrument(level = "trace", skip_all)]
470    fn reconfigure(
471        &self,
472        new: &ArtiCombinedConfig,
473        how: Reconfigure,
474    ) -> Result<(), ReconfigureError> {
475        if how == Reconfigure::AllOrNothing {
476            // If we're in all-or-nothing mode, we check it first.
477            self.reconfigure(new, Reconfigure::CheckAllOrNothing)?;
478        }
479        let dry_run = how == Reconfigure::CheckAllOrNothing;
480
481        if new.0.application().defer_bootstrap && !self.orig_defer_bootstrap {
482            how.cannot_change_specific("defer_bootstrap", "from off to on")?;
483        }
484        if !dry_run && !new.0.application().defer_bootstrap {
485            self.ensure_bootstrap_launched()
486                .map_err(into_internal!("Unable to launch client bootstrap"))?;
487        }
488
489        TorClient::reconfigure(&self.client, &new.1, how).map_err(extract_reconfigure_error)?;
490        Ok(())
491    }
492}
493
494/// If possible, extract the ReconfigureError from `err`.  Otherwise,
495/// return `err` as an internal ReconfigureError.
496//
497// (We could get rid of this function if arti_client::Error were not opaque,
498// or if arti_client::reconfigure were to return a ReconfigureError.
499// But  now is not the time to revisit those decisions.)
500fn extract_reconfigure_error(err: arti_client::Error) -> ReconfigureError {
501    for e in ErrorSources::new(&err) {
502        if let Some(reconfig_error) = e.downcast_ref::<ReconfigureError>() {
503            return reconfig_error.clone();
504        };
505    }
506    (into_internal!("Failure while reconfiguring")(err)).into()
507}
508
509impl<R: Runtime> LaunchableTorClient<R> {
510    /// Create a new LaunchableTorClient.
511    ///
512    /// We assume that it has (or has not) been told to bootstrap itself based on `cfg`.
513    pub(crate) fn new(client: Arc<TorClient<R>>, cfg: &crate::ApplicationConfig) -> Self {
514        Self {
515            orig_defer_bootstrap: cfg.defer_bootstrap,
516            have_launched: Mutex::new(!cfg.defer_bootstrap),
517            client,
518        }
519    }
520
521    /// If we have not already told this LaunchableTorClient to bootstrap itself, do so.
522    fn ensure_bootstrap_launched(&self) -> Result<(), futures::task::SpawnError> {
523        let mut have_launched = self.have_launched.lock().expect("lock poisoned");
524
525        if *have_launched {
526            return Ok(());
527        }
528
529        let client = Arc::clone(&self.client);
530        // We spawn this as a new task since `bootstrap` is very much async,
531        // but this needs to be called from `reconfigure`, which is not.
532        self.client.runtime().spawn(async move {
533            let _outcome = client.bootstrap().await;
534        })?;
535
536        *have_launched = true;
537        Ok(())
538    }
539
540    /// As [`TorClient::bootstrap`], but performs necessary bookkeeping to remember
541    /// that we have launched a bootstrap attempt.
542    pub(crate) async fn bootstrap(&self) -> arti_client::Result<()> {
543        *self.have_launched.lock().expect("lock poisoned") = true;
544
545        self.client.bootstrap().await
546    }
547}
548
549/// Internal type to represent the Arti application as a `ReconfigurableModule`.
550pub(crate) struct Application {
551    /// The configuration that Arti had at startup.
552    ///
553    /// We use this to check whether the user is asking for any impermissible
554    /// transitions.
555    original_config: ArtiConfig,
556}
557
558impl Application {
559    /// Construct a new `Application` to receive configuration changes for the
560    /// arti application.
561    pub(crate) fn new(cfg: ArtiConfig) -> Self {
562        Self {
563            original_config: cfg,
564        }
565    }
566}
567
568impl ReconfigurableModule for Application {
569    #[instrument(level = "trace", skip_all)]
570    fn reconfigure(
571        &self,
572        new: &ArtiCombinedConfig,
573        how: Reconfigure,
574    ) -> Result<(), ReconfigureError> {
575        if how == Reconfigure::AllOrNothing {
576            // If we're in all-or-nothing mode, we check it first.
577            self.reconfigure(new, Reconfigure::CheckAllOrNothing)?;
578        }
579        let dry_run = how == Reconfigure::CheckAllOrNothing;
580
581        let original = &self.original_config;
582        let config = &new.0;
583
584        if config.proxy() != original.proxy() {
585            how.cannot_change("proxy settings")?;
586        }
587        if config.logging() != original.logging() {
588            how.cannot_change("logging")?;
589        }
590        #[cfg(feature = "rpc")]
591        if config.rpc != original.rpc {
592            how.cannot_change("RPC settings")?;
593        }
594        if config.application().permit_debugging && !original.application().permit_debugging {
595            how.cannot_change_specific("application hardening", "from on to off")?;
596        }
597        // Note that this is the only config transition we actually perform so far.
598        if !dry_run && !config.application().permit_debugging {
599            #[cfg(feature = "harden")]
600            crate::process::enable_process_hardening()
601                .map_err(into_internal!("can't disable debugging"))?;
602        }
603
604        Ok(())
605    }
606}
607
608/// Find the configuration files and prepare the watcher
609fn prepare<'a, R: Runtime>(
610    watcher: &mut FileWatcherBuilder<R>,
611    sources: &'a ConfigurationSources,
612) -> anyhow::Result<FoundConfigFiles<'a>> {
613    let sources = sources.scan()?;
614    for source in sources.iter() {
615        match source {
616            ConfigurationSource::Dir(dir) => watcher.watch_dir(dir, "toml")?,
617            ConfigurationSource::File(file) => watcher.watch_path(file)?,
618            ConfigurationSource::Verbatim(_) => {}
619        }
620    }
621    Ok(sources)
622}
623
624/// Reload the configuration files, apply the runtime configuration, and
625/// reconfigure the client as much as we can.
626///
627/// Return true if we should be watching for configuration changes.
628#[instrument(level = "trace", skip_all)]
629fn reconfigure(
630    config: &ConfigurationTree,
631    mgr_inner: &mut CfgMgrInner,
632    how: Reconfigure,
633) -> Result<bool, ChangeConfigurationError> {
634    #[allow(unused_mut)]
635    let mut resolve_options = ConfigResolveOptions::default();
636    #[cfg(feature = "rpc")]
637    {
638        resolve_options.want_output_tree = true;
639    }
640
641    let rs = tor_config::resolve_return_results::<ArtiCombinedConfig>(config, &resolve_options)?;
642    let config = rs.value;
643
644    // Filter out the modules that have been dropped
645    let reconfigurable: Vec<_> = mgr_inner.modules.iter().flat_map(Weak::upgrade).collect();
646    let has_modules = !reconfigurable.is_empty();
647
648    if how == Reconfigure::AllOrNothing {
649        for module in &reconfigurable {
650            module.reconfigure(&config, Reconfigure::CheckAllOrNothing)?;
651        }
652    }
653    for module in &reconfigurable {
654        module.reconfigure(&config, how)?;
655    }
656
657    #[cfg(feature = "rpc")]
658    if how != Reconfigure::CheckAllOrNothing {
659        mgr_inner.normalized_cfg = rs
660            .output_tree
661            .expect("normalized cfg not exposed as expected!?");
662        mgr_inner.deprecated_keys = rs.deprecated.into_iter().collect();
663        mgr_inner.unrecognized_keys = rs.unrecognized.into_iter().collect();
664    }
665
666    Ok(has_modules && config.0.application().watch_configuration)
667}
668
669/// An error that occurred while trying to reload and/or replace our configuration
670#[cfg_attr(feature = "experimental-api", visibility::make(pub))]
671#[non_exhaustive]
672#[derive(thiserror::Error, Clone, Debug)]
673pub(crate) enum ChangeConfigurationError {
674    /// When we tried to make an application-level request in the configuration tree, we
675    /// were unable to do so.
676    #[error("Unable to modify configuration tree: {0}")]
677    Apply(String),
678
679    /// When we tried to merge the RPC tree into the loaded configuration, we weren't able
680    /// to do so.
681    #[error("Internal: RPC configuration tree did not apply cleanly.")]
682    Merge(#[from] tor_config::ConfigError),
683
684    /// We couldn't turn the configuration tree into the appropriate set of data structures.
685    #[error("Invalid configuration")]
686    Resolve(#[from] tor_config::load::ConfigResolveError),
687
688    /// One of the transitions we tried to make was not allowed, or failed as we tried to apply it.
689    #[error("Configuration transition failed")]
690    Transition(#[from] ReconfigureError),
691}
692
693impl HasKind for ChangeConfigurationError {
694    fn kind(&self) -> tor_error::ErrorKind {
695        use ChangeConfigurationError as E;
696        use tor_error::ErrorKind as EK;
697        match self {
698            E::Apply(_) => EK::InvalidConfig,
699            E::Merge(_) => EK::InvalidConfig,
700            E::Resolve(_) => EK::InvalidConfig,
701            E::Transition(e) => e.kind(),
702        }
703    }
704}
705
706#[cfg(test)]
707mod test {
708    // @@ begin test lint list maintained by maint/add_warning @@
709    #![allow(clippy::bool_assert_comparison)]
710    #![allow(clippy::clone_on_copy)]
711    #![allow(clippy::dbg_macro)]
712    #![allow(clippy::mixed_attributes_style)]
713    #![allow(clippy::print_stderr)]
714    #![allow(clippy::print_stdout)]
715    #![allow(clippy::single_char_pattern)]
716    #![allow(clippy::unwrap_used)]
717    #![allow(clippy::unchecked_time_subtraction)]
718    #![allow(clippy::useless_vec)]
719    #![allow(clippy::needless_pass_by_value)]
720    #![allow(clippy::string_slice)] // See arti#2571
721    //! <!-- @@ end test lint list maintained by maint/add_warning @@ -->
722
723    use crate::ArtiConfigBuilder;
724
725    use super::*;
726    use futures::SinkExt as _;
727    use futures::channel::mpsc;
728    use postage::watch;
729    use std::path::PathBuf;
730    use std::sync::{Arc, Mutex};
731    use test_temp_dir::{TestTempDir, test_temp_dir};
732    use tor_async_utils::PostageWatchSenderExt;
733    use tor_config::sources::MustRead;
734
735    /// Filename for config1
736    const CONFIG_NAME1: &str = "config1.toml";
737    /// Filename for config2
738    const CONFIG_NAME2: &str = "config2.toml";
739    /// Filename for config3
740    const CONFIG_NAME3: &str = "config3.toml";
741
742    struct TestModule {
743        // A sender for sending the new config to the test function
744        tx: Arc<Mutex<watch::Sender<ArtiCombinedConfig>>>,
745    }
746
747    impl ReconfigurableModule for TestModule {
748        fn reconfigure(
749            &self,
750            new: &ArtiCombinedConfig,
751            _how: Reconfigure,
752        ) -> Result<(), ReconfigureError> {
753            let config = new.clone();
754            self.tx.lock().unwrap().maybe_send(|_| config);
755
756            Ok(())
757        }
758    }
759
760    /// Create a test reconfigurable module.
761    ///
762    /// Returns the module and a channel on which the new configs received by the module are sent.
763    async fn create_module() -> (
764        Arc<dyn ReconfigurableModule>,
765        watch::Receiver<ArtiCombinedConfig>,
766    ) {
767        let (tx, mut rx) = watch::channel();
768        // Read the initial value from the postage::watch stream
769        // (the first observed value on this test stream is always the default config)
770        let _: ArtiCombinedConfig = rx.next().await.unwrap();
771
772        (
773            Arc::new(TestModule {
774                tx: Arc::new(Mutex::new(tx)),
775            }),
776            rx,
777        )
778    }
779
780    /// Write `data` to file `name` within `dir`.
781    fn write_file(dir: &TestTempDir, name: &str, data: &[u8]) -> PathBuf {
782        let tmp = dir.as_path_untracked().join("tmp");
783        std::fs::write(&tmp, data).unwrap();
784        let path = dir.as_path_untracked().join(name);
785        // Atomically write the config file
786        std::fs::rename(tmp, &path).unwrap();
787        path
788    }
789
790    /// Write an `ArtiConfigBuilder` to a file within `dir`.
791    fn write_config(dir: &TestTempDir, name: &str, config: &ArtiConfigBuilder) -> PathBuf {
792        let s = toml::to_string(&config).unwrap();
793        write_file(dir, name, s.as_bytes())
794    }
795
796    #[test]
797    fn watch_single_file() {
798        tor_rtcompat::test_with_one_runtime!(|rt| async move {
799            let temp_dir = test_temp_dir!();
800            let mut config_builder = ArtiConfigBuilder::default();
801            config_builder.application().watch_configuration(true);
802
803            let cfg_file = write_config(&temp_dir, CONFIG_NAME1, &config_builder);
804            let mut cfg_sources = ConfigurationSources::new_empty();
805            cfg_sources.push_source(ConfigurationSource::File(cfg_file), MustRead::MustRead);
806
807            let (module, mut rx) = create_module().await;
808
809            config_builder.logging().log_sensitive_information(true);
810            let _: PathBuf = write_config(&temp_dir, CONFIG_NAME1, &config_builder);
811
812            let (fw_tx, fw_rx) = file_watcher::channel();
813            let mgr = Arc::new(CfgMgr {
814                runtime: rt.clone(),
815                sources: cfg_sources,
816                tx: fw_tx,
817                inner: Mutex::new(CfgMgrInner {
818                    modules: vec![Arc::downgrade(&module)],
819                    ..Default::default()
820                }),
821            });
822
823            let (watcher, _) = mgr.launch_file_watcher().unwrap();
824            mgr.inner.lock().unwrap().watcher = Some(watcher);
825            let weak_mgr = Arc::downgrade(&mgr);
826
827            // Use a fake sighup stream to wait until run_watcher()'s select_biased!
828            // loop is entered
829            let (mut sighup_tx, sighup_rx) = mpsc::unbounded();
830            let runtime = rt.clone();
831            let () = rt
832                .spawn(async move {
833                    run_watcher(runtime.clone(), fw_rx, sighup_rx, weak_mgr, None)
834                        .await
835                        .unwrap();
836                })
837                .unwrap();
838
839            sighup_tx.send(()).await.unwrap();
840
841            // The reconfigurable modules should've been reloaded in response to sighup
842            let config = rx.next().await.unwrap();
843            assert_eq!(config.0, config_builder.build().unwrap());
844
845            // Overwrite the config
846            config_builder.logging().log_sensitive_information(false);
847            let _: PathBuf = write_config(&temp_dir, CONFIG_NAME1, &config_builder);
848            // The reconfigurable modules should've been reloaded in response to the config change
849            let config = rx.next().await.unwrap();
850            assert_eq!(config.0, config_builder.build().unwrap());
851        });
852    }
853
854    // TODO: Ignored until #1607 is fixed
855    #[test]
856    #[ignore]
857    fn watch_multiple() {
858        tor_rtcompat::test_with_one_runtime!(|rt| async move {
859            let temp_dir = test_temp_dir!();
860            let mut config_builder1 = ArtiConfigBuilder::default();
861            config_builder1.application().watch_configuration(true);
862
863            let _: PathBuf = write_config(&temp_dir, CONFIG_NAME1, &config_builder1);
864            let mut cfg_sources = ConfigurationSources::new_empty();
865            cfg_sources.push_source(
866                ConfigurationSource::Dir(temp_dir.as_path_untracked().to_path_buf()),
867                MustRead::MustRead,
868            );
869
870            let (module, mut rx) = create_module().await;
871
872            let (fw_tx, fw_rx) = file_watcher::channel();
873            let mgr = Arc::new(CfgMgr {
874                runtime: rt.clone(),
875                sources: cfg_sources,
876                tx: fw_tx,
877                inner: Mutex::new(CfgMgrInner {
878                    modules: vec![Arc::downgrade(&module)],
879                    ..Default::default()
880                }),
881            });
882
883            let (watcher, _) = mgr.launch_file_watcher().unwrap();
884            mgr.inner.lock().unwrap().watcher = Some(watcher);
885            let weak_mgr = Arc::downgrade(&mgr);
886
887            // Use a fake sighup stream to wait until run_watcher()'s select_biased!
888            // loop is entered
889            let (mut sighup_tx, sighup_rx) = mpsc::unbounded();
890            let runtime = rt.clone();
891            let () = rt
892                .spawn(async move {
893                    run_watcher(runtime.clone(), fw_rx, sighup_rx, weak_mgr, None)
894                        .await
895                        .unwrap();
896                })
897                .unwrap();
898
899            config_builder1.logging().log_sensitive_information(true);
900            let _: PathBuf = write_config(&temp_dir, CONFIG_NAME1, &config_builder1);
901            sighup_tx.send(()).await.unwrap();
902            // The reconfigurable modules should've been reloaded in response to sighup
903            let config = rx.next().await.unwrap();
904            assert_eq!(config.0, config_builder1.build().unwrap());
905
906            let mut config_builder2 = ArtiConfigBuilder::default();
907            config_builder2.application().watch_configuration(true);
908            // Write another config file...
909            config_builder2.system().max_files(0_u64);
910            let _: PathBuf = write_config(&temp_dir, CONFIG_NAME2, &config_builder2);
911            // Check that the 2 config files are merged
912            let mut config_builder_combined = config_builder1.clone();
913            config_builder_combined.system().max_files(0_u64);
914            let config = rx.next().await.unwrap();
915            assert_eq!(config.0, config_builder_combined.build().unwrap());
916            // Now write a new config file to the watched dir
917            config_builder2.logging().console("foo".to_string());
918            let mut config_builder_combined2 = config_builder_combined.clone();
919            config_builder_combined2
920                .logging()
921                .console("foo".to_string());
922            let config3: PathBuf = write_config(&temp_dir, CONFIG_NAME3, &config_builder2);
923            let config = rx.next().await.unwrap();
924            assert_eq!(config.0, config_builder_combined2.build().unwrap());
925
926            // Removing the file should also trigger an event
927            std::fs::remove_file(config3).unwrap();
928            let config = rx.next().await.unwrap();
929            assert_eq!(config.0, config_builder_combined.build().unwrap());
930        });
931    }
932}