use std::collections::HashMap;
use std::panic::{catch_unwind, AssertUnwindSafe};
use std::sync::{Arc, RwLock};
use crate::breadcrumb_buffer::Breadcrumb;
use crate::configuration::Configuration;
use crate::delivery_queue::DeliveryQueue;
use crate::event_builder::{capture_backtrace, Event, EventBuilder, Frame};
use crate::pii_scrubber::Value;
pub struct Reporter {
configuration: Arc<RwLock<Configuration>>,
delivery_queue: DeliveryQueue,
}
impl Reporter {
pub fn new(configuration: Arc<RwLock<Configuration>>, delivery_queue: DeliveryQueue) -> Self {
Reporter {
configuration,
delivery_queue,
}
}
#[allow(clippy::too_many_arguments)]
pub fn report(
&self,
exception_class: &str,
message: &str,
backtrace: Vec<Frame>,
context: HashMap<String, Value>,
user: Option<HashMap<String, Value>>,
breadcrumbs: Vec<Breadcrumb>,
) {
let _ = catch_unwind(AssertUnwindSafe(|| {
let config = self.configuration.read().unwrap();
if !config.is_enabled() {
return;
}
let event: Event = EventBuilder::new(&config).build(
exception_class,
message,
backtrace,
context,
user,
breadcrumbs,
);
drop(config);
self.delivery_queue.push(event);
}));
}
pub fn report_with_captured_backtrace(
&self,
exception_class: &str,
message: &str,
context: HashMap<String, Value>,
user: Option<HashMap<String, Value>>,
breadcrumbs: Vec<Breadcrumb>,
) {
let backtrace = { capture_backtrace(&self.configuration.read().unwrap()) };
self.report(
exception_class,
message,
backtrace,
context,
user,
breadcrumbs,
);
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::client::Client;
use std::net::TcpListener;
use std::sync::atomic::{AtomicI32, Ordering};
use std::time::Duration;
fn serving_reporter(status_ok: bool, environment: &str) -> (Reporter, Arc<AtomicI32>) {
use std::io::Write;
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, Ordering::SeqCst);
let status_line = if status_ok {
"HTTP/1.1 200 OK"
} else {
"HTTP/1.1 500 Internal Server Error"
};
let _ = stream.write_all(
format!("{status_line}\r\nContent-Length: 2\r\nConnection: close\r\n\r\n{{}}")
.as_bytes(),
);
}
});
let mut config = Configuration::new();
config.dsn = Some(format!("http://key@{addr}/events"));
config.environment = environment.to_string();
config.timeout = Duration::from_secs(2);
let configuration = Arc::new(RwLock::new(config));
let client = Arc::new(Client::new(Arc::clone(&configuration)));
let delivery_queue = DeliveryQueue::new(10, client);
(Reporter::new(configuration, delivery_queue), received)
}
#[test]
fn report_skips_when_disabled() {
let (reporter, received) = serving_reporter(true, "development");
reporter.report_with_captured_backtrace("Error", "boom", HashMap::new(), None, Vec::new());
std::thread::sleep(Duration::from_millis(100));
assert_eq!(received.load(Ordering::SeqCst), 0);
}
#[test]
fn report_delivers_when_enabled() {
let (reporter, received) = serving_reporter(true, "production");
reporter.report_with_captured_backtrace("Error", "boom", HashMap::new(), None, Vec::new());
let deadline = std::time::Instant::now() + Duration::from_secs(2);
while received.load(Ordering::SeqCst) < 1 && std::time::Instant::now() < deadline {
std::thread::sleep(Duration::from_millis(5));
}
assert_eq!(received.load(Ordering::SeqCst), 1);
}
#[test]
fn report_includes_the_given_user_never_scrubbed_even_though_its_an_email() {
use std::io::Write;
let listener = TcpListener::bind("127.0.0.1:0").unwrap();
let addr = listener.local_addr().unwrap();
let (tx, rx) = std::sync::mpsc::channel();
std::thread::spawn(move || {
let (mut stream, _) = listener.accept().unwrap();
let total = crate::test_support::read_full_request(&mut stream);
let _ = tx.send(String::from_utf8_lossy(&total).into_owned());
let _ = stream
.write_all(b"HTTP/1.1 200 OK\r\nContent-Length: 2\r\nConnection: close\r\n\r\n{}");
});
let mut config = Configuration::new();
config.dsn = Some(format!("http://key@{addr}/events"));
config.environment = "production".to_string();
config.timeout = Duration::from_secs(2);
let configuration = Arc::new(RwLock::new(config));
let client = Arc::new(Client::new(Arc::clone(&configuration)));
let delivery_queue = DeliveryQueue::new(10, client);
let reporter = Reporter::new(configuration, delivery_queue);
let mut user = HashMap::new();
user.insert("id".to_string(), Value::Number(42.0));
user.insert(
"email".to_string(),
Value::String("alice@example.com".to_string()),
);
reporter.report_with_captured_backtrace(
"Error",
"boom",
HashMap::new(),
Some(user),
Vec::new(),
);
let body = rx
.recv_timeout(Duration::from_secs(2))
.expect("server never received a request");
assert!(body.contains("alice@example.com"));
}
#[test]
fn report_never_panics_even_when_delivery_fails() {
let (reporter, _received) = serving_reporter(false, "production");
reporter.report_with_captured_backtrace("Error", "boom", HashMap::new(), None, Vec::new());
}
}