mod console;
mod logs;
pub mod metrics;
pub mod traces;
use std::net::ToSocketAddrs;
use std::sync::LazyLock;
use anyhow::{Result, anyhow};
use opentelemetry::global;
use opentelemetry_sdk::Resource;
use tracing::{Level, Subscriber};
use tracing_appender::non_blocking::{NonBlockingBuilder, WorkerGuard};
use tracing_subscriber::EnvFilter;
use tracing_subscriber::filter::{LevelFilter, ParseError};
use tracing_subscriber::prelude::*;
use crate::cli::LogFormat;
use crate::cli::validator::parser::tracing::CustomFilter;
use crate::cnf::ENABLE_TOKIO_CONSOLE;
pub static OTEL_DEFAULT_RESOURCE: LazyLock<Resource> = LazyLock::new(|| {
Resource::builder().with_service_name("surrealdb").build()
});
#[derive(Debug, Clone)]
pub struct Builder {
format: LogFormat,
filter: CustomFilter,
socket: Option<String>,
file_filter: Option<CustomFilter>,
otel_filter: Option<CustomFilter>,
socket_filter: Option<CustomFilter>,
socket_format: LogFormat,
file_enabled: bool,
file_format: LogFormat,
file_path: Option<String>,
file_name: Option<String>,
file_rotation: Option<String>,
}
pub fn builder() -> Builder {
Builder::default()
}
impl Default for Builder {
fn default() -> Self {
Self {
filter: CustomFilter {
env: EnvFilter::default(),
spans: std::collections::HashMap::new(),
},
format: LogFormat::Text,
socket: None,
file_filter: None,
otel_filter: None,
socket_filter: None,
socket_format: LogFormat::Text,
file_format: LogFormat::Text,
file_enabled: false,
file_path: Some("logs".to_string()),
file_name: Some("surrealdb.log".to_string()),
file_rotation: Some("daily".to_string()),
}
}
}
impl Builder {
pub fn init(self) -> Result<Vec<WorkerGuard>> {
let (registry, guards) = self.build()?;
registry.init();
Ok(guards)
}
pub fn with_filter(mut self, filter: CustomFilter) -> Self {
self.filter = filter;
self
}
pub fn with_log_level(mut self, log_level: &str) -> Self {
if let Ok(filter) = filter_from_value(log_level) {
self.filter = CustomFilter {
env: filter,
spans: std::collections::HashMap::new(),
};
}
self
}
pub fn with_file_filter(mut self, filter: Option<CustomFilter>) -> Self {
self.file_filter = filter;
self
}
pub fn with_otel_filter(mut self, filter: Option<CustomFilter>) -> Self {
self.otel_filter = filter;
self
}
pub fn with_socket_filter(mut self, filter: Option<CustomFilter>) -> Self {
self.socket_filter = filter;
self
}
pub fn with_socket(mut self, socket: Option<String>) -> Self {
self.socket = socket;
self
}
pub fn with_log_format(mut self, format: LogFormat) -> Self {
self.format = format;
self
}
pub fn with_file_format(mut self, format: LogFormat) -> Self {
self.file_format = format;
self
}
pub fn with_socket_format(mut self, format: LogFormat) -> Self {
self.format = format;
self
}
pub fn with_file_enabled(mut self, enabled: bool) -> Self {
self.file_enabled = enabled;
self
}
pub fn with_file_path(mut self, path: Option<String>) -> Self {
self.file_path = path;
self
}
pub fn with_file_name(mut self, name: Option<String>) -> Self {
self.file_name = name;
self
}
pub fn with_file_rotation(mut self, rotation: Option<String>) -> Self {
self.file_rotation = rotation;
self
}
pub fn build(&self) -> Result<(Box<dyn Subscriber + Send + Sync + 'static>, Vec<WorkerGuard>)> {
if let Some(provider) = metrics::init()? {
global::set_meter_provider(provider);
}
let (stdout, stdout_guard) = NonBlockingBuilder::default()
.lossy(true)
.thread_name("surrealdb-logger-stdout")
.finish(std::io::stdout());
let (stderr, stderr_guard) = NonBlockingBuilder::default()
.lossy(true)
.thread_name("surrealdb-logger-stderr")
.finish(std::io::stderr());
let stdio_layer = logs::output(self.filter.clone(), stdout, stderr, self.format)?;
let registry = tracing_subscriber::registry();
let registry = registry.with(stdio_layer);
let mut guards = vec![stdout_guard, stderr_guard];
let mut layers = Vec::new();
{
let filter = self.otel_filter.clone().unwrap_or_else(|| self.filter.clone());
if let Some(layer) = traces::new(filter)? {
layers.push(layer);
}
}
if let Some(addr) = &self.socket {
let address =
addr.to_socket_addrs()?.next().ok_or_else(|| anyhow!("No matching addresses"))?;
let socket = logs::socket::connect(address)?;
let (writer, guard) = NonBlockingBuilder::default()
.lossy(false)
.thread_name("surrealdb-logger-socket")
.finish(socket);
let filter = self.socket_filter.clone().unwrap_or_else(|| self.filter.clone());
let layer = logs::file(filter, writer, self.socket_format)?;
layers.push(layer);
guards.push(guard);
}
if self.file_enabled {
let file_appender = {
let path = self.file_path.as_deref().unwrap_or("logs");
let name = self.file_name.as_deref().unwrap_or("surrealdb.log");
match self.file_rotation.as_deref() {
Some("hourly") => tracing_appender::rolling::hourly(path, name),
Some("daily") => tracing_appender::rolling::daily(path, name),
Some("never") => tracing_appender::rolling::never(path, name),
_ => tracing_appender::rolling::daily(path, name),
}
};
let (writer, guard) = NonBlockingBuilder::default()
.lossy(false)
.thread_name("surrealdb-logger-file")
.finish(file_appender);
let filter = self.file_filter.clone().unwrap_or_else(|| self.filter.clone());
let layer = logs::file(filter, writer, self.file_format)?;
layers.push(layer);
guards.push(guard);
}
if *ENABLE_TOKIO_CONSOLE {
let layer = console::new()?;
layers.push(layer);
}
match layers.len() {
0 => {
Ok((Box::new(registry), guards))
}
_ => {
let registry = registry.with(layers);
Ok((Box::new(registry), guards))
}
}
}
}
pub fn shutdown() {
trace!("Shutting down telemetry service");
}
pub fn filter_from_value(v: &str) -> std::result::Result<EnvFilter, ParseError> {
match v {
"none" => Ok(EnvFilter::default()),
"error" => Ok(EnvFilter::default().add_directive(Level::ERROR.into())),
"warn" => Ok(EnvFilter::default().add_directive(Level::WARN.into())),
"info" => Ok(EnvFilter::default().add_directive(Level::INFO.into())),
"debug" => Ok(EnvFilter::default()
.add_directive(Level::WARN.into())
.add_directive("surreal=debug".parse()?)
.add_directive("surrealdb=debug".parse()?)
.add_directive("surrealdb::core::kvs::tx=debug".parse()?)
.add_directive("surrealdb::core::kvs::tr=debug".parse()?)),
"trace" => Ok(EnvFilter::default()
.add_directive(Level::WARN.into())
.add_directive("surreal=trace".parse()?)
.add_directive("surrealdb=trace".parse()?)
.add_directive("surrealdb::core::kvs::tx=debug".parse()?)
.add_directive("surrealdb::core::kvs::tr=debug".parse()?)),
"full" => Ok(EnvFilter::default()
.add_directive(Level::DEBUG.into())
.add_directive("surreal=trace".parse()?)
.add_directive("surrealdb=trace".parse()?)
.add_directive("surrealdb::core::kvs::tx=trace".parse()?)
.add_directive("surrealdb::core::kvs::tr=trace".parse()?)),
"all" => Ok(EnvFilter::default().add_directive(Level::TRACE.into())),
_ => EnvFilter::builder().parse(v),
}
}
pub fn span_filters_from_value(v: &str) -> Vec<(String, LevelFilter)> {
v.split(',')
.filter_map(|d| {
let d = d.trim();
if !d.starts_with('[') {
return None;
}
let close = d.find(']')?;
let name = &d[1..close];
let level = d[close + 1..].trim();
let level = if let Some(stripped) = level.strip_prefix('=') {
stripped.parse().ok()?
} else {
LevelFilter::TRACE
};
Some((name.to_string(), level))
})
.collect()
}