saluki_app/logging/
mod.rs1use 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
31const NB_LOG_WRITER_BUFFER_SIZE: usize = 4096;
37
38type OutputStack = Vec<Box<dyn Layer<Registry> + Send + Sync>>;
39
40pub struct LoggingGuard {
47 worker_guards: Vec<WorkerGuard>,
48 stack_handle: reload::Handle<OutputStack, Registry>,
49 controller: LoggingOverrideController,
50}
51
52impl LoggingGuard {
53 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 let _old_guards = std::mem::replace(&mut self.worker_guards, new_guards);
81
82 Ok(())
83 }
84
85 pub fn controller(&self) -> LoggingOverrideController {
87 self.controller.clone()
88 }
89}
90
91pub(crate) async fn initialize_logging(
103 config: LoggingConfiguration,
104) -> Result<(LoggingGuard, Supervisor), GenericError> {
105 let (output_stack, worker_guards) = build_output_stack(&config)?;
108 let (output_layer, stack_handle) = reload::Layer::new(output_stack);
109
110 let (filter_layer, filter_handle) = reload::Layer::new(config.log_level.as_targets());
112
113 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 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 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 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 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}