mod breadcrumb_buffer;
mod client;
mod configuration;
mod delivery_queue;
mod event_builder;
mod histogram_bucketer;
mod metric_buffer;
mod performance_flusher;
mod pii_scrubber;
mod reporter;
mod span_buffer;
mod span_queue;
mod sql_statement;
#[cfg(test)]
mod test_support;
pub use breadcrumb_buffer::Breadcrumb;
pub use configuration::Configuration;
pub use event_builder::{Event, Frame};
pub use pii_scrubber::Value;
pub use sql_statement::SqlObjects;
use std::cell::RefCell;
use std::collections::HashMap;
use std::sync::{Arc, OnceLock, RwLock};
use client::Client;
use delivery_queue::DeliveryQueue;
use metric_buffer::{Kind, MetricBuffer};
use performance_flusher::PerformanceFlusher;
use reporter::Reporter;
use span_queue::SpanQueue;
struct State {
configuration: Arc<RwLock<Configuration>>,
reporter: Reporter,
performance_flusher: Arc<PerformanceFlusher>,
span_queue: SpanQueue,
metrics: Arc<MetricBuffer>,
infrastructure_metrics: Arc<MetricBuffer>,
}
static STATE: OnceLock<State> = OnceLock::new();
thread_local! {
static CURRENT_USER: RefCell<Option<HashMap<String, Value>>> = const { RefCell::new(None) };
}
fn current_user() -> Option<HashMap<String, Value>> {
CURRENT_USER.with(|u| u.borrow().clone())
}
fn state() -> &'static State {
STATE.get_or_init(|| {
let configuration = Arc::new(RwLock::new(Configuration::new()));
let client = Arc::new(Client::new(Arc::clone(&configuration)));
let queue_size = configuration.read().unwrap().queue_size;
let delivery_queue = DeliveryQueue::new(queue_size, Arc::clone(&client));
let span_queue = SpanQueue::new(queue_size, Arc::clone(&client));
let metrics = MetricBuffer::new(
Kind::Custom,
Arc::clone(&configuration),
Arc::clone(&client),
);
let infrastructure_metrics = MetricBuffer::new(
Kind::Infrastructure,
Arc::clone(&configuration),
Arc::clone(&client),
);
let reporter = Reporter::new(Arc::clone(&configuration), delivery_queue);
let performance_flusher = PerformanceFlusher::new(Arc::clone(&configuration), client);
State {
configuration,
reporter,
performance_flusher,
span_queue,
metrics,
infrastructure_metrics,
}
})
}
pub fn init(configure: impl FnOnce(&mut Configuration)) {
let s = state();
let install_hook = {
let mut config = s.configuration.write().unwrap();
configure(&mut config);
config.install_panic_hook
};
if install_hook {
install_panic_hook();
}
}
pub fn capture_error<E: std::error::Error>(
err: &E,
context: HashMap<String, Value>,
user: Option<HashMap<String, Value>>,
) {
capture_error_with_class(std::any::type_name::<E>(), err, context, user);
}
pub fn capture_error_with_class(
exception_class: &str,
err: &dyn std::error::Error,
context: HashMap<String, Value>,
user: Option<HashMap<String, Value>>,
) {
let s = state();
s.reporter.report_with_captured_backtrace(
exception_class,
&err.to_string(),
context,
user.or_else(current_user),
breadcrumb_buffer::current_breadcrumbs(),
);
}
pub fn capture_error_with_sql<E: std::error::Error>(
err: &E,
statement: &str,
context: HashMap<String, Value>,
user: Option<HashMap<String, Value>>,
) {
capture_error_with_class_and_sql(std::any::type_name::<E>(), err, statement, context, user);
}
pub fn capture_error_with_class_and_sql(
exception_class: &str,
err: &dyn std::error::Error,
statement: &str,
context: HashMap<String, Value>,
user: Option<HashMap<String, Value>>,
) {
let s = state();
s.reporter.report_with_captured_backtrace_and_sql(
exception_class,
&err.to_string(),
context,
user.or_else(current_user),
breadcrumb_buffer::current_breadcrumbs(),
statement,
);
}
pub fn set_user(user: HashMap<String, Value>) {
let user = if user.is_empty() { None } else { Some(user) };
CURRENT_USER.with(|u| *u.borrow_mut() = user);
}
pub fn add_breadcrumb(message: &str, category: &str, level: &str, data: HashMap<String, Value>) {
let config = state().configuration.read().unwrap();
breadcrumb_buffer::add_breadcrumb(
config.track_breadcrumbs,
config.max_breadcrumbs,
message,
category,
level,
data,
);
}
pub fn clear_breadcrumbs() {
breadcrumb_buffer::clear_breadcrumbs();
}
pub fn record_performance(transaction_name: &str, duration_ms: f64) {
state()
.performance_flusher
.record(transaction_name, duration_ms);
}
pub fn time_transaction<T>(transaction_name: &str, f: impl FnOnce() -> T) -> T {
struct Timing<'a> {
name: &'a str,
started_at: std::time::Instant,
}
impl Drop for Timing<'_> {
fn drop(&mut self) {
record_performance(self.name, self.started_at.elapsed().as_secs_f64() * 1000.0);
}
}
let _timing = Timing {
name: transaction_name,
started_at: std::time::Instant::now(),
};
f()
}
pub fn flush_performance() {
state().performance_flusher.flush();
}
pub fn capture_metric(name: &str, value: f64) {
let s = state();
let (enabled, environment, release) = {
let config = s.configuration.read().unwrap();
(
config.is_enabled(),
config.environment.clone(),
config.release.clone(),
)
};
if !enabled {
return;
}
let release = release
.as_deref()
.map(pii_scrubber::json_string)
.unwrap_or_else(|| "null".to_string());
s.metrics.record(
value,
&format!(
"\"metric_name\":{},\"value\":{value},\"environment\":{},\"release\":{release}",
pii_scrubber::json_string(name),
pii_scrubber::json_string(&environment),
),
);
}
pub fn capture_infrastructure_metric(name: &str, value: f64, hostname: Option<&str>) {
let s = state();
let (enabled, server_name) = {
let config = s.configuration.read().unwrap();
(config.is_enabled(), config.server_name.clone())
};
if !enabled {
return;
}
let hostname = hostname
.map(str::to_string)
.or(server_name)
.unwrap_or_default();
s.infrastructure_metrics.record(
value,
&format!(
"\"metric_name\":{},\"value\":{value},\"hostname\":{}",
pii_scrubber::json_string(name),
pii_scrubber::json_string(&hostname),
),
);
}
pub fn flush_metrics() {
let s = state();
s.metrics.flush();
s.infrastructure_metrics.flush();
}
pub fn trace<T>(root_name: &str, f: impl FnOnce() -> T) -> T {
let (owns_trace, threshold) = {
let config = state().configuration.read().unwrap();
(span_buffer::begin(&config), config.trace_capture_threshold)
};
if !owns_trace {
return span(root_name, "service", HashMap::new(), f);
}
struct Root<'a> {
name: &'a str,
started_at: std::time::SystemTime,
timer: std::time::Instant,
threshold: std::time::Duration,
}
impl Drop for Root<'_> {
fn drop(&mut self) {
let duration_ms = self.timer.elapsed().as_secs_f64() * 1000.0;
if let Some(body) = span_buffer::end(
self.threshold,
self.name,
"controller",
self.started_at,
duration_ms,
) {
state().span_queue.push(body);
}
}
}
let _root = Root {
name: root_name,
started_at: std::time::SystemTime::now(),
timer: std::time::Instant::now(),
threshold,
};
f()
}
pub fn span<T>(name: &str, kind: &str, data: HashMap<String, Value>, f: impl FnOnce() -> T) -> T {
let Some((id, parent)) = span_buffer::open_span() else {
return f();
};
struct Open<'a> {
id: String,
parent: String,
name: &'a str,
kind: &'a str,
data: HashMap<String, Value>,
started_at: std::time::SystemTime,
timer: std::time::Instant,
}
impl Drop for Open<'_> {
fn drop(&mut self) {
span_buffer::close_span(
std::mem::take(&mut self.id),
std::mem::take(&mut self.parent),
self.name,
self.kind,
self.started_at,
self.timer.elapsed().as_secs_f64() * 1000.0,
std::mem::take(&mut self.data),
);
}
}
let _open = Open {
id,
parent,
name,
kind,
data,
started_at: std::time::SystemTime::now(),
timer: std::time::Instant::now(),
};
f()
}
pub fn record_span(
name: &str,
kind: &str,
started_at: std::time::SystemTime,
duration_ms: f64,
data: HashMap<String, Value>,
) {
span_buffer::record_leaf(name, kind, started_at, duration_ms, data);
}
pub fn install_panic_hook() {
let previous = std::panic::take_hook();
std::panic::set_hook(Box::new(move |info| {
report_panic(info);
previous(info);
}));
}
fn report_panic(info: &std::panic::PanicHookInfo) {
let s = state();
let (exception_class, message) = panic_message(info);
let mut context = HashMap::new();
if let Some(location) = info.location() {
context.insert(
"panic_location".to_string(),
Value::String(format!(
"{}:{}:{}",
location.file(),
location.line(),
location.column()
)),
);
}
s.reporter.report_with_captured_backtrace(
&exception_class,
&message,
context,
current_user(),
breadcrumb_buffer::current_breadcrumbs(),
);
}
fn panic_message(info: &std::panic::PanicHookInfo) -> (String, String) {
let payload = info.payload();
if let Some(s) = payload.downcast_ref::<&str>() {
("panic".to_string(), s.to_string())
} else if let Some(s) = payload.downcast_ref::<String>() {
("panic".to_string(), s.clone())
} else {
("panic".to_string(), "non-string panic payload".to_string())
}
}
#[macro_export]
macro_rules! context {
( $( $key:expr => $value:expr ),* $(,)? ) => {{
#[allow(unused_mut)]
let mut map = ::std::collections::HashMap::new();
$( map.insert(::std::string::ToString::to_string($key), $crate::Value::from($value)); )*
map
}};
}
pub trait ResultReportExt<T> {
fn report_err(
self,
context: HashMap<String, Value>,
user: Option<HashMap<String, Value>>,
) -> Self;
}
impl<T, E: std::error::Error> ResultReportExt<T> for Result<T, E> {
fn report_err(
self,
context: HashMap<String, Value>,
user: Option<HashMap<String, Value>>,
) -> Self {
if let Err(ref err) = self {
capture_error(err, context, user);
}
self
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::io::{Read, Write};
use std::net::TcpListener;
use std::sync::atomic::{AtomicBool, AtomicI32, Ordering as AtomicOrdering};
use std::sync::Mutex;
use std::time::Duration;
fn spawn_tracker_server() -> (std::net::SocketAddr, Arc<AtomicI32>) {
let listener = TcpListener::bind("127.0.0.1:0").unwrap();
let addr = listener.local_addr().unwrap();
let received = Arc::new(AtomicI32::new(0));
let received_clone = Arc::clone(&received);
std::thread::spawn(move || {
for stream in listener.incoming().flatten() {
let mut stream = stream;
crate::test_support::read_full_request(&mut stream);
received_clone.fetch_add(1, AtomicOrdering::SeqCst);
let _ = stream.write_all(
b"HTTP/1.1 200 OK\r\nContent-Length: 2\r\nConnection: close\r\n\r\n{}",
);
}
});
(addr, received)
}
fn spawn_tracker_server_capturing_body(
) -> (std::net::SocketAddr, Arc<AtomicI32>, Arc<Mutex<String>>) {
let listener = TcpListener::bind("127.0.0.1:0").unwrap();
let addr = listener.local_addr().unwrap();
let received = Arc::new(AtomicI32::new(0));
let received_clone = Arc::clone(&received);
let last_body = Arc::new(Mutex::new(String::new()));
let last_body_clone = Arc::clone(&last_body);
std::thread::spawn(move || {
for stream in listener.incoming().flatten() {
let mut stream = stream;
let mut buf = [0u8; 8192];
let mut total = Vec::new();
let mut header_end = None;
loop {
let n = stream.read(&mut buf).unwrap_or(0);
if n == 0 {
break;
}
total.extend_from_slice(&buf[..n]);
if let Some(pos) = total.windows(4).position(|w| w == b"\r\n\r\n") {
header_end = Some(pos + 4);
let header_text = String::from_utf8_lossy(&total[..pos]).into_owned();
let content_length: usize = header_text
.lines()
.find(|l| l.to_lowercase().starts_with("content-length:"))
.and_then(|l| l.split(':').nth(1))
.and_then(|v| v.trim().parse().ok())
.unwrap_or(0);
while total.len() < pos + 4 + content_length {
let n = stream.read(&mut buf).unwrap_or(0);
if n == 0 {
break;
}
total.extend_from_slice(&buf[..n]);
}
break;
}
}
if let Some(header_end) = header_end {
let body = String::from_utf8_lossy(&total[header_end..]).into_owned();
*last_body_clone.lock().unwrap() = body;
}
received_clone.fetch_add(1, AtomicOrdering::SeqCst);
let _ = stream.write_all(
b"HTTP/1.1 200 OK\r\nContent-Length: 2\r\nConnection: close\r\n\r\n{}",
);
}
});
(addr, received, last_body)
}
fn wait_for(received: &AtomicI32, count: i32) {
let deadline = std::time::Instant::now() + Duration::from_secs(2);
while received.load(AtomicOrdering::SeqCst) < count && std::time::Instant::now() < deadline
{
std::thread::sleep(Duration::from_millis(5));
}
assert_eq!(received.load(AtomicOrdering::SeqCst), count);
}
#[test]
fn public_api_end_to_end() {
let (addr, received) = spawn_tracker_server();
init(|c| {
c.dsn = Some(format!("http://key@{addr}/events"));
c.environment = "production".to_string();
c.timeout = Duration::from_secs(2);
c.install_panic_hook = false; });
let err = std::io::Error::other("boom");
capture_error(&err, context! {"order_id" => 7}, None);
wait_for(&received, 1);
let (addr, received) = spawn_tracker_server();
init(|c| c.dsn = Some(format!("http://key@{addr}/events")));
let boxed: Box<dyn std::error::Error> = Box::new(std::io::Error::other("boxed boom"));
capture_error_with_class("std::io::Error", boxed.as_ref(), HashMap::new(), None);
wait_for(&received, 1);
let (addr, received) = spawn_tracker_server();
init(|c| c.dsn = Some(format!("http://key@{addr}/events")));
let result: Result<(), std::io::Error> = Err(std::io::Error::other("reported via ext"));
let passed_through = result.report_err(HashMap::new(), None);
assert!(passed_through.is_err());
wait_for(&received, 1);
let ok: Result<i32, std::io::Error> = Ok(42);
assert_eq!(ok.report_err(HashMap::new(), None).unwrap(), 42);
let (addr, received, last_body) = spawn_tracker_server_capturing_body();
init(|c| c.dsn = Some(format!("http://key@{addr}/events")));
set_user(context! {"id" => 42, "email" => "alice@example.com"});
capture_error(&std::io::Error::other("boom"), HashMap::new(), None);
wait_for(&received, 1);
assert!(last_body.lock().unwrap().contains("alice@example.com"));
let (addr, received, last_body) = spawn_tracker_server_capturing_body();
init(|c| c.dsn = Some(format!("http://key@{addr}/events")));
set_user(context! {"id" => 42});
capture_error(
&std::io::Error::other("boom"),
HashMap::new(),
Some(context! {"id" => 99}),
);
wait_for(&received, 1);
assert!(last_body.lock().unwrap().contains("\"id\":99"));
set_user(HashMap::new());
let (addr, received, last_body) = spawn_tracker_server_capturing_body();
init(|c| c.dsn = Some(format!("http://key@{addr}/events")));
add_breadcrumb("opened checkout", "custom", "info", HashMap::new());
add_breadcrumb(
"charged card",
"custom",
"info",
context! {"order_id" => 42},
);
capture_error(&std::io::Error::other("boom"), HashMap::new(), None);
wait_for(&received, 1);
{
let body = last_body.lock().unwrap();
assert!(body.contains("opened checkout"));
assert!(body.contains("charged card"));
}
clear_breadcrumbs();
let (addr, received, last_body) = spawn_tracker_server_capturing_body();
init(|c| c.dsn = Some(format!("http://key@{addr}/events")));
capture_error(&std::io::Error::other("boom"), HashMap::new(), None);
wait_for(&received, 1);
assert!(!last_body.lock().unwrap().contains("\"breadcrumbs\""));
let (addr, received, last_body) = spawn_tracker_server_capturing_body();
init(|c| c.dsn = Some(format!("http://key@{addr}/events")));
record_performance("GET /users/:id", 10.0);
record_performance("GET /users/:id", 30.0);
flush_performance();
wait_for(&received, 1);
{
let body = last_body.lock().unwrap();
assert!(body.starts_with("{\"samples\":["), "body = {body}");
assert!(body.contains("\"transaction_name\":\"GET /users/:id\""));
assert!(body.contains("\"request_count\":2"));
}
let (addr, received, last_body) = spawn_tracker_server_capturing_body();
init(|c| c.dsn = Some(format!("http://key@{addr}/events")));
let value = time_transaction("timed", || 42);
assert_eq!(value, 42);
let panicked = std::panic::catch_unwind(|| {
time_transaction("timed panics", || -> u32 {
panic!("inside a timed transaction")
})
});
assert!(panicked.is_err());
flush_performance();
wait_for(&received, 1);
{
let body = last_body.lock().unwrap();
assert!(
body.contains("\"transaction_name\":\"timed\""),
"body = {body}"
);
assert!(
body.contains("\"transaction_name\":\"timed panics\""),
"body = {body}"
);
}
let (addr, received, _last_body) = spawn_tracker_server_capturing_body();
init(|c| {
c.dsn = Some(format!("http://key@{addr}/events"));
c.track_performance = false;
});
record_performance("never recorded", 10.0);
flush_performance();
std::thread::sleep(Duration::from_millis(150));
assert_eq!(received.load(AtomicOrdering::SeqCst), 0);
init(|c| c.track_performance = true);
let (addr, received) = spawn_tracker_server();
let previous_hook_ran = Arc::new(AtomicBool::new(false));
let previous_hook_ran_clone = Arc::clone(&previous_hook_ran);
std::panic::set_hook(Box::new(move |_| {
previous_hook_ran_clone.store(true, AtomicOrdering::SeqCst);
}));
init(|c| {
c.dsn = Some(format!("http://key@{addr}/events"));
c.install_panic_hook = true;
});
let result = std::panic::catch_unwind(|| {
panic!("test panic");
});
assert!(result.is_err());
assert!(
previous_hook_ran.load(AtomicOrdering::SeqCst),
"the previously-installed hook should still have run"
);
wait_for(&received, 1);
std::panic::set_hook(Box::new(|_| {}));
}
#[test]
fn span_just_runs_the_closure_outside_a_trace() {
assert_eq!(span("free", "service", HashMap::new(), || 7), 7);
record_span(
"free",
"database",
std::time::SystemTime::now(),
1.0,
HashMap::new(),
);
}
#[test]
fn a_panic_inside_trace_or_span_propagates_and_leaves_no_trace_behind() {
let result = std::panic::catch_unwind(|| {
trace("GET /boom", || {
span("bad", "service", HashMap::new(), || panic!("boom"));
})
});
assert!(result.is_err());
assert!(!span_buffer::is_active());
}
#[test]
fn trace_returns_the_closures_value_and_a_nested_trace_is_a_span() {
let value = trace("outer", || trace("inner", || 5));
assert_eq!(value, 5);
assert!(!span_buffer::is_active());
}
#[test]
fn context_macro_builds_expected_map() {
let ctx = context! {"order_id" => 42, "customer" => "acme-inc"};
assert_eq!(ctx.get("order_id"), Some(&Value::Number(42.0)));
assert_eq!(
ctx.get("customer"),
Some(&Value::String("acme-inc".to_string()))
);
}
}