pub mod audit_logs;
mod console;
mod logs;
pub mod metrics;
pub mod traces;
use std::net::ToSocketAddrs;
use std::sync::{LazyLock, OnceLock};
use anyhow::{Result, anyhow};
use opentelemetry::KeyValue;
use opentelemetry_sdk::Resource;
use opentelemetry_sdk::metrics::SdkMeterProvider;
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;
use crate::observe::ObservabilityRuntime;
pub static OTEL_DEFAULT_RESOURCE: LazyLock<Resource> = LazyLock::new(|| {
let edition = SERVICE_EDITION.get().copied().unwrap_or(DEFAULT_SERVICE_EDITION);
Resource::builder()
.with_service_name("surrealdb")
.with_attribute(KeyValue::new("service.edition", edition))
.build()
});
const DEFAULT_SERVICE_EDITION: &str = "community";
static SERVICE_EDITION: OnceLock<&'static str> = OnceLock::new();
pub fn set_service_edition(edition: &'static str) {
let _ = SERVICE_EDITION.set(edition);
}
pub fn service_edition() -> &'static str {
SERVICE_EDITION.get().copied().unwrap_or(DEFAULT_SERVICE_EDITION)
}
#[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()
}
fn warn_removed_env_vars() {
if std::env::var("SURREAL_TELEMETRY_NAMESPACE").ok().filter(|v| !v.trim().is_empty()).is_some()
{
warn!(
"SURREAL_TELEMETRY_NAMESPACE is set but is no longer applied. \
The `namespace` attribute was removed from telemetry metrics \
because it is tenant-identifying in multi-tenant deployments. \
Remove this variable from your deployment configuration."
);
}
if std::env::var("SURREAL_TELEMETRY_RPC_LIVE_ID")
.ok()
.filter(|v| !v.trim().is_empty())
.is_some()
{
warn!(
"SURREAL_TELEMETRY_RPC_LIVE_ID is set but is no longer applied. \
Per-notification OTLP attribution by `rpc.live_id` was removed \
when WebSocket telemetry was unified into the ExecutionObserver \
pipeline. Notification volume is still surfaced via the \
`surrealdb_live_query_notifications_total` Prometheus counter. \
Remove this variable from your deployment configuration."
);
}
}
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()),
}
}
}
pub struct TelemetryHandles {
pub guards: Vec<WorkerGuard>,
pub runtime: ObservabilityRuntime,
}
impl Builder {
pub fn init(self) -> Result<TelemetryHandles> {
let (registry, handles) = self.build()?;
registry.init();
warn_removed_env_vars();
Ok(handles)
}
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.socket_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>, TelemetryHandles)> {
let metrics_init = metrics::init()?;
let audit_logger_provider = audit_logs::init()?;
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_layers = logs::output(self.filter.clone(), stdout, stderr, self.format)?;
let registry = tracing_subscriber::registry();
let registry = registry.with(stdio_layers);
let mut guards = vec![stdout_guard, stderr_guard];
let mut layers = Vec::new();
let mut tracer_provider = None;
{
let filter = self.otel_filter.clone().unwrap_or_else(default_otel_filter);
if let Some(trace_layer) = traces::new(filter)? {
layers.push(trace_layer.layer);
tracer_provider = Some(trace_layer.provider);
opentelemetry::global::set_text_map_propagator(
opentelemetry_sdk::propagation::TraceContextPropagator::new(),
);
}
}
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);
}
let mut runtime_builder = match metrics_init {
Some(init) => {
let mut b = ObservabilityRuntime::builder(init.provider)
.with_resource(OTEL_DEFAULT_RESOURCE.clone());
if let Some(exporter) = init.prometheus_exporter {
b = b.with_prometheus_exporter(exporter);
}
b
}
None => ObservabilityRuntime::builder(SdkMeterProvider::default())
.with_resource(OTEL_DEFAULT_RESOURCE.clone()),
};
if let Some(provider) = audit_logger_provider {
runtime_builder = runtime_builder.with_audit_logger_provider(provider);
}
if let Some(provider) = tracer_provider {
runtime_builder = runtime_builder.with_tracer_provider(provider);
}
let handles = TelemetryHandles {
guards,
runtime: runtime_builder.build(),
};
match layers.len() {
0 => {
Ok((Box::new(registry), handles))
}
_ => {
let registry = registry.with(layers);
Ok((Box::new(registry), handles))
}
}
}
}
pub fn shutdown() {
trace!("Shutting down telemetry service");
}
fn default_otel_filter() -> CustomFilter {
let env = EnvFilter::default()
.add_directive(Level::INFO.into())
.add_directive("surrealdb::core::rpc=debug".parse().expect("static filter directive"))
.add_directive("surrealdb::core::dbs=debug".parse().expect("static filter directive"))
.add_directive(
"surrealdb_server::ntw::tracer=debug".parse().expect("static filter directive"),
)
.add_directive(
"surrealdb_server::telemetry::traces::rpc=debug"
.parse()
.expect("static filter directive"),
);
CustomFilter {
env,
spans: std::collections::HashMap::new(),
}
}
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()
}