Skip to main content

common_telemetry/
logging.rs

1// Copyright 2023 Greptime Team
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15//! logging stuffs, inspired by databend
16mod file_retention;
17
18use std::collections::HashMap;
19use std::env;
20use std::io::IsTerminal;
21use std::sync::{Arc, Mutex, Once, RwLock};
22use std::time::Duration;
23
24use common_base::readable_size::ReadableSize;
25use common_base::serde::empty_string_as_default;
26use file_retention::{DirectoryRetention, LogFileKind, build_file_appender};
27use once_cell::sync::{Lazy, OnceCell};
28use opentelemetry::trace::TracerProvider;
29use opentelemetry::{KeyValue, global};
30use opentelemetry_otlp::{Protocol, SpanExporter, WithExportConfig, WithHttpConfig};
31use opentelemetry_sdk::propagation::TraceContextPropagator;
32use opentelemetry_sdk::trace::{Sampler, Tracer};
33use opentelemetry_semantic_conventions::resource;
34use serde::{Deserialize, Serialize};
35use tracing::callsite;
36use tracing::metadata::LevelFilter;
37use tracing_appender::non_blocking::WorkerGuard;
38use tracing_log::LogTracer;
39use tracing_subscriber::filter::{FilterFn, Targets};
40use tracing_subscriber::fmt::Layer;
41use tracing_subscriber::layer::{Layered, SubscriberExt};
42use tracing_subscriber::prelude::*;
43use tracing_subscriber::{EnvFilter, Registry, filter};
44
45use crate::tracing_sampler::{TracingSampleOptions, create_sampler};
46
47/// The default endpoint when use gRPC exporter protocol.
48pub const DEFAULT_OTLP_GRPC_ENDPOINT: &str = "http://localhost:4317";
49
50/// The default endpoint when use HTTP exporter protocol.
51pub const DEFAULT_OTLP_HTTP_ENDPOINT: &str = "http://localhost:4318/v1/traces";
52
53/// The default logs directory.
54pub const DEFAULT_LOGGING_DIR: &str = "logs";
55
56/// Handle for reloading log level
57pub static LOG_RELOAD_HANDLE: OnceCell<tracing_subscriber::reload::Handle<Targets, Registry>> =
58    OnceCell::new();
59
60type DynSubscriber = Layered<tracing_subscriber::reload::Layer<Targets, Registry>, Registry>;
61type OtelTraceLayer = tracing_opentelemetry::OpenTelemetryLayer<DynSubscriber, Tracer>;
62
63#[derive(Clone)]
64pub struct TraceReloadHandle {
65    inner: Arc<RwLock<Option<OtelTraceLayer>>>,
66}
67
68impl TraceReloadHandle {
69    fn new(inner: Arc<RwLock<Option<OtelTraceLayer>>>) -> Self {
70        Self { inner }
71    }
72
73    pub fn reload(&self, new_layer: Option<OtelTraceLayer>) {
74        let mut guard = self.inner.write().unwrap();
75        *guard = new_layer;
76        drop(guard);
77
78        callsite::rebuild_interest_cache();
79    }
80}
81
82/// A tracing layer that can be dynamically reloaded.
83///
84/// Mostly copied from [`tracing_subscriber::reload::Layer`].
85struct TraceLayer {
86    inner: Arc<RwLock<Option<OtelTraceLayer>>>,
87}
88
89impl TraceLayer {
90    fn new(initial: Option<OtelTraceLayer>) -> (Self, TraceReloadHandle) {
91        let inner = Arc::new(RwLock::new(initial));
92        (
93            Self {
94                inner: inner.clone(),
95            },
96            TraceReloadHandle::new(inner),
97        )
98    }
99
100    fn with_layer<R>(&self, f: impl FnOnce(&OtelTraceLayer) -> R) -> Option<R> {
101        self.inner
102            .read()
103            .ok()
104            .and_then(|guard| guard.as_ref().map(f))
105    }
106
107    fn with_layer_mut<R>(&self, f: impl FnOnce(&mut OtelTraceLayer) -> R) -> Option<R> {
108        self.inner
109            .write()
110            .ok()
111            .and_then(|mut guard| guard.as_mut().map(f))
112    }
113}
114
115impl tracing_subscriber::Layer<DynSubscriber> for TraceLayer {
116    fn on_register_dispatch(&self, subscriber: &tracing::Dispatch) {
117        let _ = self.with_layer(|layer| layer.on_register_dispatch(subscriber));
118    }
119
120    fn on_layer(&mut self, subscriber: &mut DynSubscriber) {
121        let _ = self.with_layer_mut(|layer| layer.on_layer(subscriber));
122    }
123
124    fn register_callsite(
125        &self,
126        metadata: &'static tracing::Metadata<'static>,
127    ) -> tracing::subscriber::Interest {
128        self.with_layer(|layer| layer.register_callsite(metadata))
129            .unwrap_or_else(tracing::subscriber::Interest::always)
130    }
131
132    fn enabled(
133        &self,
134        metadata: &tracing::Metadata<'_>,
135        ctx: tracing_subscriber::layer::Context<'_, DynSubscriber>,
136    ) -> bool {
137        self.with_layer(|layer| layer.enabled(metadata, ctx))
138            .unwrap_or(true)
139    }
140
141    fn on_new_span(
142        &self,
143        attrs: &tracing::span::Attributes<'_>,
144        id: &tracing::span::Id,
145        ctx: tracing_subscriber::layer::Context<'_, DynSubscriber>,
146    ) {
147        let _ = self.with_layer(|layer| layer.on_new_span(attrs, id, ctx));
148    }
149
150    fn max_level_hint(&self) -> Option<LevelFilter> {
151        self.with_layer(|layer| layer.max_level_hint()).flatten()
152    }
153
154    fn on_record(
155        &self,
156        span: &tracing::span::Id,
157        values: &tracing::span::Record<'_>,
158        ctx: tracing_subscriber::layer::Context<'_, DynSubscriber>,
159    ) {
160        let _ = self.with_layer(|layer| layer.on_record(span, values, ctx));
161    }
162
163    fn on_follows_from(
164        &self,
165        span: &tracing::span::Id,
166        follows: &tracing::span::Id,
167        ctx: tracing_subscriber::layer::Context<'_, DynSubscriber>,
168    ) {
169        let _ = self.with_layer(|layer| layer.on_follows_from(span, follows, ctx));
170    }
171
172    fn event_enabled(
173        &self,
174        event: &tracing::Event<'_>,
175        ctx: tracing_subscriber::layer::Context<'_, DynSubscriber>,
176    ) -> bool {
177        self.with_layer(|layer| layer.event_enabled(event, ctx))
178            .unwrap_or(true)
179    }
180
181    fn on_event(
182        &self,
183        event: &tracing::Event<'_>,
184        ctx: tracing_subscriber::layer::Context<'_, DynSubscriber>,
185    ) {
186        let _ = self.with_layer(|layer| layer.on_event(event, ctx));
187    }
188
189    fn on_enter(
190        &self,
191        id: &tracing::span::Id,
192        ctx: tracing_subscriber::layer::Context<'_, DynSubscriber>,
193    ) {
194        let _ = self.with_layer(|layer| layer.on_enter(id, ctx));
195    }
196
197    fn on_exit(
198        &self,
199        id: &tracing::span::Id,
200        ctx: tracing_subscriber::layer::Context<'_, DynSubscriber>,
201    ) {
202        let _ = self.with_layer(|layer| layer.on_exit(id, ctx));
203    }
204
205    fn on_close(
206        &self,
207        id: tracing::span::Id,
208        ctx: tracing_subscriber::layer::Context<'_, DynSubscriber>,
209    ) {
210        let _ = self.with_layer(|layer| layer.on_close(id, ctx));
211    }
212
213    fn on_id_change(
214        &self,
215        old: &tracing::span::Id,
216        new: &tracing::span::Id,
217        ctx: tracing_subscriber::layer::Context<'_, DynSubscriber>,
218    ) {
219        let _ = self.with_layer(|layer| layer.on_id_change(old, new, ctx));
220    }
221
222    unsafe fn downcast_raw(&self, id: std::any::TypeId) -> Option<*const ()> {
223        self.inner.read().ok().and_then(|guard| {
224            guard
225                .as_ref()
226                .and_then(|layer| unsafe { layer.downcast_raw(id) })
227        })
228    }
229}
230
231/// Handle for reloading trace level
232pub static TRACE_RELOAD_HANDLE: OnceCell<TraceReloadHandle> = OnceCell::new();
233
234static TRACER: OnceCell<Mutex<TraceState>> = OnceCell::new();
235
236#[derive(Debug)]
237enum TraceState {
238    Ready(Tracer),
239    Deferred(TraceContext),
240}
241
242/// The logging options that used to initialize the logger.
243#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
244#[serde(default)]
245pub struct LoggingOptions {
246    /// The directory to store log files. If not set, logs will be written to stdout.
247    pub dir: String,
248
249    /// The log level that can be one of "trace", "debug", "info", "warn", "error". Default is "info".
250    pub level: Option<String>,
251
252    /// The log format that can be one of "json" or "text". Default is "text".
253    #[serde(default, deserialize_with = "empty_string_as_default")]
254    pub log_format: LogFormat,
255
256    /// The maximum number of log files set by default.
257    pub max_log_files: usize,
258
259    /// The maximum total size of managed log files in `dir`. Zero disables size-based retention.
260    pub max_log_dir_size: ReadableSize,
261
262    /// Whether to append logs to stdout. Default is true.
263    pub append_stdout: bool,
264
265    /// Whether to write logs to files in `dir`. Default is true.
266    pub enable_file_logging: bool,
267
268    /// Whether to enable tracing with OTLP. Default is false.
269    pub enable_otlp_tracing: bool,
270
271    /// The endpoint of OTLP.
272    pub otlp_endpoint: Option<String>,
273
274    /// The tracing sample ratio.
275    pub tracing_sample_ratio: Option<TracingSampleOptions>,
276
277    /// The protocol of OTLP export.
278    pub otlp_export_protocol: Option<OtlpExportProtocol>,
279
280    /// Additional HTTP headers for OTLP exporter.
281    #[serde(skip_serializing_if = "HashMap::is_empty")]
282    pub otlp_headers: HashMap<String, String>,
283
284    /// Whether to enable per-region metrics.
285    pub enable_per_region_metrics: bool,
286}
287
288/// The protocol of OTLP export.
289#[derive(Clone, Debug, Serialize, Deserialize, PartialEq)]
290#[serde(rename_all = "snake_case")]
291pub enum OtlpExportProtocol {
292    /// GRPC protocol.
293    Grpc,
294
295    /// HTTP protocol with binary protobuf.
296    Http,
297}
298
299/// The options of slow query.
300#[derive(Clone, Debug, Serialize, Deserialize, PartialEq)]
301#[serde(default)]
302pub struct SlowQueryOptions {
303    /// Whether to enable slow query log.
304    pub enable: bool,
305
306    /// The record type of slow queries.
307    #[serde(deserialize_with = "empty_string_as_default")]
308    pub record_type: SlowQueriesRecordType,
309
310    /// The threshold of slow queries.
311    #[serde(with = "humantime_serde")]
312    pub threshold: Duration,
313
314    /// The sample ratio of slow queries.
315    pub sample_ratio: f64,
316
317    /// The table TTL of `slow_queries` system table. Default is "90d".
318    /// It's used when `record_type` is `SystemTable`.
319    #[serde(with = "humantime_serde")]
320    pub ttl: Duration,
321}
322
323impl Default for SlowQueryOptions {
324    fn default() -> Self {
325        Self {
326            enable: true,
327            record_type: SlowQueriesRecordType::SystemTable,
328            threshold: Duration::from_secs(30),
329            sample_ratio: 1.0,
330            ttl: Duration::from_days(90),
331        }
332    }
333}
334
335#[derive(Clone, Debug, Serialize, Deserialize, Copy, PartialEq, Default)]
336#[serde(rename_all = "snake_case")]
337pub enum SlowQueriesRecordType {
338    /// Record the slow query in the system table.
339    #[default]
340    SystemTable,
341    /// Record the slow query in a specific logs file.
342    Log,
343}
344
345#[derive(Clone, Debug, Copy, PartialEq, Eq, Serialize, Deserialize, Default)]
346#[serde(rename_all = "snake_case")]
347pub enum LogFormat {
348    Json,
349    #[default]
350    Text,
351}
352
353#[derive(Clone, Debug)]
354struct TraceContext {
355    app_name: String,
356    node_id: String,
357    logging_opts: LoggingOptions,
358}
359
360impl Default for LoggingOptions {
361    fn default() -> Self {
362        Self {
363            // The directory path will be configured at application startup, typically using the data home directory as a base.
364            dir: "".to_string(),
365            level: None,
366            log_format: LogFormat::Text,
367            enable_otlp_tracing: false,
368            otlp_endpoint: None,
369            tracing_sample_ratio: None,
370            append_stdout: true,
371            enable_file_logging: true,
372            // Rotation hourly, 24 files per day, keeps info log files of 30 days
373            max_log_files: 720,
374            max_log_dir_size: ReadableSize::default(),
375            otlp_export_protocol: None,
376            otlp_headers: HashMap::new(),
377            enable_per_region_metrics: false,
378        }
379    }
380}
381
382#[derive(Default, Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
383pub struct TracingOptions {
384    #[cfg(feature = "tokio-console")]
385    pub tokio_console_addr: Option<String>,
386}
387
388/// Init tracing for unittest.
389/// Write logs to file `unittest`.
390pub fn init_default_ut_logging() {
391    static START: Once = Once::new();
392
393    START.call_once(|| {
394        let mut g = GLOBAL_UT_LOG_GUARD.as_ref().lock().unwrap();
395
396        // When running in Github's actions, env "UNITTEST_LOG_DIR" is set to a directory other
397        // than "/tmp".
398        // This is to fix the problem that the "/tmp" disk space of action runner's is small,
399        // if we write testing logs in it, actions would fail due to disk out of space error.
400        let dir =
401            env::var("UNITTEST_LOG_DIR").unwrap_or_else(|_| "/tmp/__unittest_logs".to_string());
402
403        let level = env::var("UNITTEST_LOG_LEVEL").unwrap_or_else(|_|
404            "debug,hyper=warn,tower=warn,datafusion=warn,reqwest=warn,sqlparser=warn,h2=info,opendal=info,rskafka=info".to_string()
405        );
406        let opts = LoggingOptions {
407            dir: dir.clone(),
408            level: Some(level),
409            ..Default::default()
410        };
411        *g = Some(init_global_logging(
412            "unittest",
413            &opts,
414            &TracingOptions::default(),
415            None,
416            None,
417        ));
418
419        crate::info!("logs dir = {}", dir);
420    });
421}
422
423static GLOBAL_UT_LOG_GUARD: Lazy<Arc<Mutex<Option<Vec<WorkerGuard>>>>> =
424    Lazy::new(|| Arc::new(Mutex::new(None)));
425
426const DEFAULT_LOG_TARGETS: &str = "info";
427
428#[allow(clippy::print_stdout)]
429pub fn init_global_logging(
430    app_name: &str,
431    opts: &LoggingOptions,
432    tracing_opts: &TracingOptions,
433    node_id: Option<String>,
434    slow_query_opts: Option<&SlowQueryOptions>,
435) -> Vec<WorkerGuard> {
436    static START: Once = Once::new();
437    let mut guards = vec![];
438    let node_id = node_id.unwrap_or_else(|| "none".to_string());
439
440    START.call_once(|| {
441        // Enable log compatible layer to convert log record to tracing span.
442        LogTracer::init().expect("log tracer must be valid");
443
444        // Configure the stdout logging layer.
445        let stdout_logging_layer = if opts.append_stdout {
446            let (writer, guard) = tracing_appender::non_blocking(std::io::stdout());
447            guards.push(guard);
448
449            if opts.log_format == LogFormat::Json {
450                Some(
451                    Layer::new()
452                        .json()
453                        .with_writer(writer)
454                        .with_ansi(std::io::stdout().is_terminal())
455                        .boxed(),
456                )
457            } else {
458                Some(
459                    Layer::new()
460                        .with_writer(writer)
461                        .with_ansi(std::io::stdout().is_terminal())
462                        .boxed(),
463                )
464            }
465        } else {
466            None
467        };
468
469        let file_logging_enabled = opts.enable_file_logging && !opts.dir.is_empty();
470
471        let retention = file_logging_enabled
472            .then(|| {
473                DirectoryRetention::new(opts.dir.clone(), opts.max_log_dir_size, opts.max_log_files)
474            })
475            .flatten();
476
477        // Configure the file logging layer with rolling policy.
478        let file_logging_layer = if file_logging_enabled {
479            let rolling_appender =
480                build_file_appender(opts, LogFileKind::Default, retention.as_ref());
481            let (writer, guard) = tracing_appender::non_blocking(rolling_appender);
482            guards.push(guard);
483
484            if opts.log_format == LogFormat::Json {
485                Some(
486                    Layer::new()
487                        .json()
488                        .with_writer(writer)
489                        .with_ansi(false)
490                        .boxed(),
491                )
492            } else {
493                Some(Layer::new().with_writer(writer).with_ansi(false).boxed())
494            }
495        } else {
496            None
497        };
498
499        // Configure the error file logging layer with rolling policy.
500        let err_file_logging_layer = if file_logging_enabled {
501            let rolling_appender =
502                build_file_appender(opts, LogFileKind::Error, retention.as_ref());
503            let (writer, guard) = tracing_appender::non_blocking(rolling_appender);
504            guards.push(guard);
505
506            if opts.log_format == LogFormat::Json {
507                Some(
508                    Layer::new()
509                        .json()
510                        .with_writer(writer)
511                        .with_ansi(false)
512                        .with_filter(filter::LevelFilter::ERROR)
513                        .boxed(),
514                )
515            } else {
516                Some(
517                    Layer::new()
518                        .with_writer(writer)
519                        .with_ansi(false)
520                        .with_filter(filter::LevelFilter::ERROR)
521                        .boxed(),
522                )
523            }
524        } else {
525            None
526        };
527
528        let slow_query_logging_layer =
529            build_slow_query_logger(opts, slow_query_opts, retention.as_ref(), &mut guards);
530
531        if let Some(retention) = &retention {
532            retention.initialize();
533        }
534
535        // resolve log level settings from:
536        // - options from command line or config files
537        // - environment variable: RUST_LOG
538        // - default settings
539        let filter = opts
540            .level
541            .as_deref()
542            .or(env::var(EnvFilter::DEFAULT_ENV).ok().as_deref())
543            .unwrap_or(DEFAULT_LOG_TARGETS)
544            .parse::<filter::Targets>()
545            .expect("error parsing log level string");
546
547        let (dyn_filter, reload_handle) = tracing_subscriber::reload::Layer::new(filter.clone());
548
549        LOG_RELOAD_HANDLE
550            .set(reload_handle)
551            .expect("reload handle already set, maybe init_global_logging get called twice?");
552
553        let mut initial_tracer = None;
554        let trace_state = if opts.enable_otlp_tracing {
555            let tracer = create_tracer(app_name, &node_id, opts);
556            initial_tracer = Some(tracer.clone());
557            TraceState::Ready(tracer)
558        } else {
559            TraceState::Deferred(TraceContext {
560                app_name: app_name.to_string(),
561                node_id: node_id.clone(),
562                logging_opts: opts.clone(),
563            })
564        };
565
566        TRACER
567            .set(Mutex::new(trace_state))
568            .expect("trace state already initialized");
569
570        let initial_trace_layer = initial_tracer
571            .as_ref()
572            .map(|tracer| tracing_opentelemetry::layer().with_tracer(tracer.clone()));
573
574        let (dyn_trace_layer, trace_reload_handle) = TraceLayer::new(initial_trace_layer);
575
576        TRACE_RELOAD_HANDLE
577            .set(trace_reload_handle)
578            .unwrap_or_else(|_| panic!("failed to set trace reload handle"));
579
580        // Must enable 'tokio_unstable' cfg to use this feature.
581        // For example: `RUSTFLAGS="--cfg tokio_unstable" cargo run -F common-telemetry/console -- standalone start`
582        #[cfg(feature = "tokio-console")]
583        let subscriber = {
584            let tokio_console_layer =
585                if let Some(tokio_console_addr) = &tracing_opts.tokio_console_addr {
586                    let addr: std::net::SocketAddr = tokio_console_addr.parse().unwrap_or_else(|e| {
587                    panic!("Invalid binding address '{tokio_console_addr}' for tokio-console: {e}");
588                });
589                    println!("tokio-console listening on {addr}");
590
591                    Some(
592                        console_subscriber::ConsoleLayer::builder()
593                            .server_addr(addr)
594                            .spawn(),
595                    )
596                } else {
597                    None
598                };
599
600            Registry::default()
601                .with(dyn_filter)
602                .with(dyn_trace_layer)
603                .with(tokio_console_layer)
604                .with(stdout_logging_layer)
605                .with(file_logging_layer)
606                .with(err_file_logging_layer)
607                .with(slow_query_logging_layer)
608        };
609
610        // consume the `tracing_opts` to avoid "unused" warnings.
611        let _ = tracing_opts;
612
613        #[cfg(not(feature = "tokio-console"))]
614        let subscriber = Registry::default()
615            .with(dyn_filter)
616            .with(dyn_trace_layer)
617            .with(stdout_logging_layer)
618            .with(file_logging_layer)
619            .with(err_file_logging_layer)
620            .with(slow_query_logging_layer);
621
622        global::set_text_map_propagator(TraceContextPropagator::new());
623
624        tracing::subscriber::set_global_default(subscriber)
625            .expect("error setting global tracing subscriber");
626    });
627
628    guards
629}
630
631fn create_tracer(app_name: &str, node_id: &str, opts: &LoggingOptions) -> Tracer {
632    let sampler = opts
633        .tracing_sample_ratio
634        .as_ref()
635        .map(create_sampler)
636        .map(Sampler::ParentBased)
637        .unwrap_or(Sampler::ParentBased(Box::new(Sampler::AlwaysOn)));
638
639    let resource = opentelemetry_sdk::Resource::builder_empty()
640        .with_attributes([
641            KeyValue::new(resource::SERVICE_NAME, app_name.to_string()),
642            KeyValue::new(resource::SERVICE_INSTANCE_ID, node_id.to_string()),
643            KeyValue::new(resource::SERVICE_VERSION, common_version::version()),
644            KeyValue::new(resource::PROCESS_PID, std::process::id().to_string()),
645        ])
646        .build();
647
648    opentelemetry_sdk::trace::SdkTracerProvider::builder()
649        .with_batch_exporter(build_otlp_exporter(opts))
650        .with_sampler(sampler)
651        .with_resource(resource)
652        .build()
653        .tracer("greptimedb")
654}
655
656/// Ensure that the OTLP tracer has been constructed, building it lazily if needed.
657pub fn get_or_init_tracer() -> Result<Tracer, &'static str> {
658    let state = TRACER.get().ok_or("trace state is not initialized")?;
659    let mut guard = state.lock().expect("trace state lock poisoned");
660
661    match &mut *guard {
662        TraceState::Ready(tracer) => Ok(tracer.clone()),
663        TraceState::Deferred(context) => {
664            let tracer = create_tracer(&context.app_name, &context.node_id, &context.logging_opts);
665            *guard = TraceState::Ready(tracer.clone());
666            Ok(tracer)
667        }
668    }
669}
670
671fn build_otlp_exporter(opts: &LoggingOptions) -> SpanExporter {
672    let protocol = opts
673        .otlp_export_protocol
674        .clone()
675        .unwrap_or(OtlpExportProtocol::Http);
676
677    let endpoint = opts
678        .otlp_endpoint
679        .as_ref()
680        .map(|e| {
681            if e.starts_with("http") {
682                e.clone()
683            } else {
684                format!("http://{}", e)
685            }
686        })
687        .unwrap_or_else(|| match protocol {
688            OtlpExportProtocol::Grpc => DEFAULT_OTLP_GRPC_ENDPOINT.to_string(),
689            OtlpExportProtocol::Http => DEFAULT_OTLP_HTTP_ENDPOINT.to_string(),
690        });
691
692    match protocol {
693        OtlpExportProtocol::Grpc => SpanExporter::builder()
694            .with_tonic()
695            .with_endpoint(endpoint)
696            .build()
697            .expect("Failed to create OTLP gRPC exporter "),
698
699        OtlpExportProtocol::Http => SpanExporter::builder()
700            .with_http()
701            .with_endpoint(endpoint)
702            .with_protocol(Protocol::HttpBinary)
703            .with_headers(opts.otlp_headers.clone())
704            .build()
705            .expect("Failed to create OTLP HTTP exporter "),
706    }
707}
708
709fn build_slow_query_logger<S>(
710    opts: &LoggingOptions,
711    slow_query_opts: Option<&SlowQueryOptions>,
712    retention: Option<&DirectoryRetention>,
713    guards: &mut Vec<WorkerGuard>,
714) -> Option<Box<dyn tracing_subscriber::Layer<S> + Send + Sync + 'static>>
715where
716    S: tracing::Subscriber
717        + Send
718        + 'static
719        + for<'span> tracing_subscriber::registry::LookupSpan<'span>,
720{
721    if let Some(slow_query_opts) = slow_query_opts {
722        if opts.enable_file_logging
723            && !opts.dir.is_empty()
724            && slow_query_opts.enable
725            && slow_query_opts.record_type == SlowQueriesRecordType::Log
726        {
727            let rolling_appender = build_file_appender(opts, LogFileKind::SlowQuery, retention);
728            let (writer, guard) = tracing_appender::non_blocking(rolling_appender);
729            guards.push(guard);
730
731            // Only logs if the field contains "slow".
732            let slow_query_filter = FilterFn::new(|metadata| {
733                metadata
734                    .fields()
735                    .iter()
736                    .any(|field| field.name().contains("slow"))
737            });
738
739            if opts.log_format == LogFormat::Json {
740                Some(
741                    Layer::new()
742                        .json()
743                        .with_writer(writer)
744                        .with_ansi(false)
745                        .with_filter(slow_query_filter)
746                        .boxed(),
747                )
748            } else {
749                Some(
750                    Layer::new()
751                        .with_writer(writer)
752                        .with_ansi(false)
753                        .with_filter(slow_query_filter)
754                        .boxed(),
755                )
756            }
757        } else {
758            None
759        }
760    } else {
761        None
762    }
763}
764
765#[cfg(test)]
766mod tests {
767    use super::*;
768
769    #[test]
770    fn test_logging_options_deserialization_default() {
771        let json = r#"{}"#;
772        let opts: LoggingOptions = serde_json::from_str(json).unwrap();
773
774        assert_eq!(opts.log_format, LogFormat::Text);
775        assert_eq!(opts.dir, "");
776        assert_eq!(opts.level, None);
777        assert!(opts.append_stdout);
778        assert!(opts.enable_file_logging);
779    }
780
781    #[test]
782    fn test_logging_options_deserialization_enable_file_logging() {
783        let json = r#"{"enable_file_logging": false}"#;
784        let opts: LoggingOptions = serde_json::from_str(json).unwrap();
785
786        assert!(!opts.enable_file_logging);
787    }
788
789    #[test]
790    fn test_logging_options_deserialization_max_log_dir_size() {
791        let json = r#"{"max_log_dir_size": "1MiB"}"#;
792        let opts: LoggingOptions = serde_json::from_str(json).unwrap();
793
794        assert_eq!(opts.max_log_dir_size, ReadableSize::mb(1));
795    }
796
797    #[test]
798    fn test_logging_options_deserialization_empty_log_format() {
799        let json = r#"{"log_format": ""}"#;
800        let opts: LoggingOptions = serde_json::from_str(json).unwrap();
801
802        // Empty string should use default (Text)
803        assert_eq!(opts.log_format, LogFormat::Text);
804    }
805
806    #[test]
807    fn test_logging_options_deserialization_valid_log_format() {
808        let json_format = r#"{"log_format": "json"}"#;
809        let opts: LoggingOptions = serde_json::from_str(json_format).unwrap();
810        assert_eq!(opts.log_format, LogFormat::Json);
811
812        let text_format = r#"{"log_format": "text"}"#;
813        let opts: LoggingOptions = serde_json::from_str(text_format).unwrap();
814        assert_eq!(opts.log_format, LogFormat::Text);
815    }
816
817    #[test]
818    fn test_logging_options_deserialization_missing_log_format() {
819        let json = r#"{"dir": "/tmp/logs"}"#;
820        let opts: LoggingOptions = serde_json::from_str(json).unwrap();
821
822        // Missing log_format should use default (Text)
823        assert_eq!(opts.log_format, LogFormat::Text);
824        assert_eq!(opts.dir, "/tmp/logs");
825    }
826
827    #[test]
828    fn test_slow_query_options_deserialization_default() {
829        let json = r#"{"enable": true, "threshold": "30s"}"#;
830        let opts: SlowQueryOptions = serde_json::from_str(json).unwrap();
831
832        assert_eq!(opts.record_type, SlowQueriesRecordType::SystemTable);
833        assert!(opts.enable);
834    }
835
836    #[test]
837    fn test_slow_query_options_deserialization_empty_record_type() {
838        let json = r#"{"enable": true, "record_type": "", "threshold": "30s"}"#;
839        let opts: SlowQueryOptions = serde_json::from_str(json).unwrap();
840
841        // Empty string should use default (SystemTable)
842        assert_eq!(opts.record_type, SlowQueriesRecordType::SystemTable);
843        assert!(opts.enable);
844    }
845
846    #[test]
847    fn test_slow_query_options_deserialization_valid_record_type() {
848        let system_table_json =
849            r#"{"enable": true, "record_type": "system_table", "threshold": "30s"}"#;
850        let opts: SlowQueryOptions = serde_json::from_str(system_table_json).unwrap();
851        assert_eq!(opts.record_type, SlowQueriesRecordType::SystemTable);
852
853        let log_json = r#"{"enable": true, "record_type": "log", "threshold": "30s"}"#;
854        let opts: SlowQueryOptions = serde_json::from_str(log_json).unwrap();
855        assert_eq!(opts.record_type, SlowQueriesRecordType::Log);
856    }
857
858    #[test]
859    fn test_slow_query_options_deserialization_missing_record_type() {
860        let json = r#"{"enable": false, "threshold": "30s"}"#;
861        let opts: SlowQueryOptions = serde_json::from_str(json).unwrap();
862
863        // Missing record_type should use default (SystemTable)
864        assert_eq!(opts.record_type, SlowQueriesRecordType::SystemTable);
865        assert!(!opts.enable);
866    }
867
868    #[test]
869    fn test_otlp_export_protocol_deserialization_valid_values() {
870        let grpc_json = r#""grpc""#;
871        let protocol: OtlpExportProtocol = serde_json::from_str(grpc_json).unwrap();
872        assert_eq!(protocol, OtlpExportProtocol::Grpc);
873
874        let http_json = r#""http""#;
875        let protocol: OtlpExportProtocol = serde_json::from_str(http_json).unwrap();
876        assert_eq!(protocol, OtlpExportProtocol::Http);
877    }
878
879    #[test]
880    fn test_logging_options_partial_eq_all_fields() {
881        let base = LoggingOptions::default();
882
883        let mut log_format = base.clone();
884        log_format.log_format = LogFormat::Json;
885        assert_ne!(base, log_format);
886
887        let mut max_log_files = base.clone();
888        max_log_files.max_log_files += 1;
889        assert_ne!(base, max_log_files);
890
891        let mut max_log_dir_size = base.clone();
892        max_log_dir_size.max_log_dir_size = ReadableSize::mb(1);
893        assert_ne!(base, max_log_dir_size);
894
895        let mut otlp_export_protocol = base.clone();
896        otlp_export_protocol.otlp_export_protocol = Some(OtlpExportProtocol::Http);
897        assert_ne!(base, otlp_export_protocol);
898
899        let mut otlp_headers = base.clone();
900        otlp_headers
901            .otlp_headers
902            .insert("key".to_string(), "value".to_string());
903        assert_ne!(base, otlp_headers);
904    }
905}