saluki_app/logging/
mod.rs

1//! Logging.
2
3// TODO: `AgentLikeFieldVisitor` currently allocates a `String` to hold the message field when it finds it. This is
4// suboptimal because it means we allocate a string literally every time we log a message. Logging is rare, but it's
5// just a recipe for small, unnecessary allocations over time... and makes it that much more inefficient to enable
6// debug/trace-level logging in production.
7//
8// We might consider _something_ like a string pool in the future, but we can defer that until we have a better idea of
9// what the potential impact is in practice.
10
11use std::io::Write;
12
13use saluki_core::runtime::Supervisor;
14use saluki_error::{generic_error, ErrorContext as _, GenericError};
15use tracing_appender::non_blocking::{NonBlocking, NonBlockingBuilder, WorkerGuard};
16use tracing_rolling_file::RollingFileAppenderBase;
17use tracing_subscriber::{layer::SubscriberExt as _, reload, util::SubscriberInitExt as _, Layer, Registry};
18
19mod api;
20pub use self::api::{LoggingAPIHandler, LoggingOverrideController, LoggingOverrideWorker};
21
22mod config;
23pub use self::config::{LogLevel, LoggingConfiguration};
24
25mod layer;
26use self::layer::{build_formatting_layer, build_syslog_formatting_layer};
27
28mod syslog;
29use self::syslog::SyslogWriter;
30
31// Number of buffered lines in each non-blocking log writer.
32//
33// This directly influences the idle memory usage since each logging backend (console, file, etc) will have a bounded
34// channel that can hold this many elements, and each element is roughly 32 bytes, so 1,000 elements/lines consumes a
35// minimum of ~32KB, etc.
36const NB_LOG_WRITER_BUFFER_SIZE: usize = 4096;
37
38type OutputStack = Vec<Box<dyn Layer<Registry> + Send + Sync>>;
39
40/// A handle to the dynamic logging subsystem.
41///
42/// Held by [`BootstrapGuard`][crate::bootstrap::BootstrapGuard] for the lifetime of the application. Owns the
43/// worker guards (which flush buffered output on drop), and exposes [`reload`][Self::reload] for swapping the
44/// entire logging configuration plus [`controller`][Self::controller] for driving runtime filter changes through
45/// the override worker.
46pub struct LoggingGuard {
47    worker_guards: Vec<WorkerGuard>,
48    stack_handle: reload::Handle<OutputStack, Registry>,
49    controller: LoggingOverrideController,
50}
51
52impl LoggingGuard {
53    /// Reloads the logging subsystem from the given configuration.
54    ///
55    /// Rebuilds the output layer stack from `config` and routes the new level filter through the override worker
56    /// as the new base filter, so an active override is preserved (the new base will take effect once the override
57    /// expires). Worker guards for the previous outputs are dropped after the swap, which flushes any buffered log
58    /// lines to their original destinations.
59    ///
60    /// This is the right entry point when the entire logging configuration may have changed (for example, outputs,
61    /// format, level). For runtime base-filter changes only -- such as following a `log_level` config update --
62    /// use [`controller`][Self::controller] and call
63    /// [`update_base`][LoggingOverrideController::update_base] directly.
64    ///
65    /// # Errors
66    ///
67    /// Returns an error if the new output layers can't be constructed (for example, the configured log file path is
68    /// inaccessible) or if the override worker is no longer running.
69    pub async fn reload(&mut self, config: LoggingConfiguration) -> Result<(), GenericError> {
70        let (new_stack, new_guards) = build_output_stack(&config)?;
71        let new_filter = config.log_level.as_targets();
72
73        self.stack_handle
74            .reload(new_stack)
75            .map_err(|e| generic_error!("Failed to swap logging output stack: {}", e))?;
76        self.controller.update_base(new_filter).await?;
77
78        // Drop the old worker guards _after_ the swap so any buffered lines are flushed to their original destinations
79        // before the worker threads exit.
80        let _old_guards = std::mem::replace(&mut self.worker_guards, new_guards);
81
82        Ok(())
83    }
84
85    /// Returns a logging override controller that can be used to change the default filter directives.
86    pub fn controller(&self) -> LoggingOverrideController {
87        self.controller.clone()
88    }
89}
90
91/// Initializes the logging subsystem for `tracing` with the ability to dynamically update the log filtering directives
92/// at runtime.
93///
94/// Returns a [`LoggingGuard`] which must be held until the application is about to shutdown, plus a [`Supervisor`]
95/// holding the subsystem's background workers. The caller must arrange for that supervisor to run -- typically by
96/// adding it to a parent supervisor -- or override requests are accepted but never applied.
97///
98/// # Errors
99///
100/// If the logging subsystem was already initialized, or its supervisor can't be constructed, an error will be
101/// returned.
102pub(crate) async fn initialize_logging(
103    config: LoggingConfiguration,
104) -> Result<(LoggingGuard, Supervisor), GenericError> {
105    // Build the initial output stack from the supplied configuration. This is later swappable via
106    // `BootstrapGuard::reload_logging` once the Datadog Agent provides authoritative configuration.
107    let (output_stack, worker_guards) = build_output_stack(&config)?;
108    let (output_layer, stack_handle) = reload::Layer::new(output_stack);
109
110    // Set up our log level filtering and dynamic filter layer.
111    let (filter_layer, filter_handle) = reload::Layer::new(config.log_level.as_targets());
112
113    // The override worker owns the canonical base filter -- the directives the system restores to after an override
114    // expires or is reset. It seeds the base from the reload handle on startup and is updated via the controller,
115    // both by `LoggingGuard::reload` once the Agent's configuration is applied and by any other caller (for example, a
116    // runtime `log_level` watcher) wired up via [`LoggingGuard::controller`].
117    let (override_worker, controller) = LoggingOverrideWorker::new(filter_handle);
118
119    tracing_subscriber::registry()
120        .with(output_layer.with_filter(filter_layer))
121        .try_init()?;
122
123    // The override worker also asserts the privileged API routes for runtime filter control, so nothing driven
124    // through those routes takes effect until this supervisor is running.
125    let mut supervisor = Supervisor::new("logging").error_context("Failed to construct logging supervisor.")?;
126    supervisor.add_worker(override_worker);
127
128    Ok((
129        LoggingGuard {
130            worker_guards,
131            stack_handle,
132            controller,
133        },
134        supervisor,
135    ))
136}
137
138fn build_output_stack(config: &LoggingConfiguration) -> Result<(OutputStack, Vec<WorkerGuard>), GenericError> {
139    let mut layers: OutputStack = Vec::new();
140    let mut guards = Vec::new();
141
142    if config.log_to_console {
143        let (nb_stdout, guard) = writer_to_nonblocking("console", std::io::stdout());
144        guards.push(guard);
145        layers.push(build_formatting_layer(config, nb_stdout));
146    }
147
148    if !config.log_file.is_empty() {
149        let appender_builder = RollingFileAppenderBase::builder();
150        let appender = appender_builder
151            .filename(config.log_file.clone())
152            .max_filecount(config.log_file_max_rolls)
153            .condition_max_file_size(config.log_file_max_size.as_u64())
154            .build()
155            .map_err(|e| generic_error!("Failed to build log file appender: {}", e))?;
156
157        let (nb_appender, guard) = writer_to_nonblocking("file", appender);
158        guards.push(guard);
159        layers.push(build_formatting_layer(config, nb_appender));
160    }
161
162    if config.log_to_syslog {
163        let syslog_writer = SyslogWriter::from_uri(&config.syslog_uri)
164            .map_err(|e| generic_error!("Failed to build syslog log writer: {}", e))?;
165        // Keep syslog on the same lossy non-blocking path as console/file so logging never
166        // backpressures ADP.
167        let (nb_syslog, guard) = writer_to_nonblocking("syslog", syslog_writer);
168        guards.push(guard);
169        layers.push(build_syslog_formatting_layer(config, nb_syslog));
170    }
171
172    Ok((layers, guards))
173}
174
175fn writer_to_nonblocking<W>(writer_name: &'static str, writer: W) -> (NonBlocking, WorkerGuard)
176where
177    W: Write + Send + 'static,
178{
179    let thread_name = format!("log-writer-{}", writer_name);
180    NonBlockingBuilder::default()
181        .thread_name(&thread_name)
182        .buffered_lines_limit(NB_LOG_WRITER_BUFFER_SIZE)
183        .lossy(true)
184        .finish(writer)
185}
186
187#[cfg(test)]
188mod tests {
189    use super::*;
190
191    const TEST_SYSLOG_URI: &str = "udp://127.0.0.1:9";
192
193    #[test]
194    fn output_stack_skips_syslog_when_disabled() {
195        let config = logging_config_without_outputs();
196
197        let (layers, guards) = build_output_stack(&config).expect("build output stack");
198
199        assert_eq!(layers.len(), 0);
200        assert_eq!(guards.len(), 0);
201    }
202
203    #[test]
204    fn output_stack_adds_syslog_layer_and_guard_when_enabled() {
205        let config = logging_config_with_syslog(TEST_SYSLOG_URI);
206
207        let (layers, guards) = build_output_stack(&config).expect("build output stack with syslog");
208
209        assert_eq!(layers.len(), 1);
210        assert_eq!(guards.len(), 1);
211    }
212
213    #[test]
214    fn output_stack_fails_when_syslog_uri_is_invalid() {
215        let config = logging_config_with_syslog("http://127.0.0.1:514");
216
217        let error = match build_output_stack(&config) {
218            Ok(_) => panic!("invalid syslog URI should fail output stack build"),
219            Err(error) => error,
220        };
221
222        let error = error.to_string();
223        assert!(error.contains("Failed to build syslog log writer"));
224        assert!(error.contains("Unsupported syslog URI scheme 'http'"));
225    }
226
227    #[tokio::test]
228    async fn reload_can_enable_change_disable_syslog_and_preserves_previous_stack_on_invalid_config() {
229        use saluki_core::runtime::Supervisor;
230        use tokio::sync::oneshot;
231
232        let config = logging_config_without_outputs();
233        let (output_stack, worker_guards) = build_output_stack(&config).expect("build initial output stack");
234        let (output_layer, stack_handle) = reload::Layer::new(output_stack);
235        let (filter_layer, filter_handle) = reload::Layer::new(config.log_level.as_targets());
236        let (override_worker, controller) = LoggingOverrideWorker::new(filter_handle);
237        let mut guard = LoggingGuard {
238            worker_guards,
239            stack_handle,
240            controller,
241        };
242        let _keep_layers_alive = (output_layer, filter_layer);
243
244        // Run the override worker inside a Supervisor so it has the dataspace context required for
245        // its route assertion. The supervisor's shutdown signal drives the worker to exit cleanly.
246        let mut sup = Supervisor::new("test-logging-override").expect("create supervisor");
247        sup.add_worker(override_worker);
248        let (shutdown_tx, shutdown_rx) = oneshot::channel::<()>();
249        let sup_handle = tokio::spawn(async move { sup.run_with_shutdown(shutdown_rx).await });
250
251        guard
252            .reload(logging_config_with_syslog(TEST_SYSLOG_URI))
253            .await
254            .expect("reload should enable syslog");
255        assert_eq!(guard.worker_guards.len(), 1);
256
257        guard
258            .reload(logging_config_with_syslog("udp://127.0.0.1:10"))
259            .await
260            .expect("reload should change syslog URI");
261        assert_eq!(guard.worker_guards.len(), 1);
262
263        guard
264            .reload(logging_config_without_outputs())
265            .await
266            .expect("reload should disable syslog");
267        assert_eq!(guard.worker_guards.len(), 0);
268
269        guard
270            .reload(logging_config_with_syslog(TEST_SYSLOG_URI))
271            .await
272            .expect("reload should re-enable syslog");
273        assert_eq!(guard.worker_guards.len(), 1);
274
275        let error = guard
276            .reload(logging_config_with_syslog("http://127.0.0.1:514"))
277            .await
278            .expect_err("invalid syslog URI should fail reload");
279        assert!(error.to_string().contains("Failed to build syslog log writer"));
280        assert_eq!(guard.worker_guards.len(), 1);
281
282        // Trigger supervisor shutdown so the worker exits via its shutdown branch rather than via
283        // the channel-close path (which would prompt a restart attempt by the supervisor).
284        shutdown_tx.send(()).expect("send shutdown");
285        sup_handle
286            .await
287            .expect("supervisor task joins")
288            .expect("supervisor should exit cleanly");
289        drop(guard);
290    }
291
292    fn logging_config_without_outputs() -> LoggingConfiguration {
293        let mut config = LoggingConfiguration::simple();
294        config.log_to_console = false;
295        config.log_file.clear();
296        config.log_to_syslog = false;
297        config.syslog_uri.clear();
298        config
299    }
300
301    fn logging_config_with_syslog(uri: &str) -> LoggingConfiguration {
302        let mut config = logging_config_without_outputs();
303        config.log_to_syslog = true;
304        config.syslog_uri = uri.to_string();
305        config
306    }
307}