introduce log_buffer type to cleanup concurrency & correctness

also cargo fmt
This commit is contained in:
Chris Beck
2025-12-05 19:40:33 -07:00
parent a214b8e474
commit f899d391b9
5 changed files with 87 additions and 49 deletions
+5 -1
View File
@@ -164,7 +164,11 @@ pub struct JsonLogMessage {
pub message: String,
/// Log level - accepts various formats (error, ERROR, err, etc.)
#[serde(default, alias = "severity", deserialize_with = "deserialize_opt_level")]
#[serde(
default,
alias = "severity",
deserialize_with = "deserialize_opt_level"
)]
pub level: Option<Level>,
/// Timestamp - accepts Unix epoch seconds (int or string) or RFC3339 string
+3 -9
View File
@@ -129,15 +129,9 @@ mod tests {
map.get_or_insert_with("b".to_string(), || 2, |_| ()).await;
map.get_or_insert_with("c".to_string(), || 3, |_| ()).await;
let a = map
.get_or_insert_with("a".to_string(), || 0, |v| *v)
.await;
let b = map
.get_or_insert_with("b".to_string(), || 0, |v| *v)
.await;
let c = map
.get_or_insert_with("c".to_string(), || 0, |v| *v)
.await;
let a = map.get_or_insert_with("a".to_string(), || 0, |v| *v).await;
let b = map.get_or_insert_with("b".to_string(), || 0, |v| *v).await;
let c = map.get_or_insert_with("c".to_string(), || 0, |v| *v).await;
assert_eq!(a, 1);
assert_eq!(b, 2);
+58
View File
@@ -0,0 +1,58 @@
//! Thread-safe log message buffer.
//!
//! Wraps a circular buffer with a synchronous mutex for fast, blocking access.
use super::circular_buffer::CircularBuffer;
use crate::log_message::LogMessage;
use std::sync::Mutex;
/// A thread-safe circular buffer for log messages.
///
/// Uses a synchronous mutex since all operations are fast and non-blocking.
#[derive(Debug)]
pub struct LogBuffer {
buf: Mutex<CircularBuffer<LogMessage>>,
}
impl LogBuffer {
/// Create a new log buffer with the given capacity.
pub fn new(capacity: usize) -> Self {
Self {
buf: Mutex::new(CircularBuffer::new(capacity)),
}
}
/// Push a log message to the buffer.
pub fn push_back(&self, msg: LogMessage) {
let mut buf = self.buf.lock().unwrap();
buf.push_back(msg);
}
/// Push a log message and drain all messages, calling `f` for each.
///
/// Messages are passed to `f` in reverse order (newest first).
/// The buffer is cleared after draining.
pub fn push_back_and_drain(&self, msg: LogMessage, mut f: impl FnMut(&LogMessage)) {
let mut buf = self.buf.lock().unwrap();
buf.push_back(msg);
for log_msg in buf.iter().rev() {
f(log_msg);
}
buf.clear();
}
/// Iterate over all messages without modifying the buffer.
///
/// Messages are passed to `f` in reverse order (newest first).
pub fn for_each(&self, mut f: impl FnMut(&LogMessage)) {
let buf = self.buf.lock().unwrap();
for log_msg in buf.iter().rev() {
f(log_msg);
}
}
/// Returns the number of messages currently in the buffer.
pub fn len(&self) -> usize {
self.buf.lock().unwrap().len()
}
}
+19 -35
View File
@@ -1,4 +1,4 @@
use super::circular_buffer::CircularBuffer;
use super::log_buffer::LogBuffer;
use super::route::{Destination, Limit, Route};
use super::{AdminMessage, LimitResult, Limiter, LimiterSet};
use crate::{
@@ -77,7 +77,7 @@ pub struct LogHandler {
config: LogHandlerConfig,
admin_mq_tx: UnboundedSender<AdminMessage>,
/// Log buffers keyed by origin (app + host). Lazily created.
log_buffers: ConcurrentMap<Origin, Mutex<CircularBuffer<LogMessage>>>,
log_buffers: ConcurrentMap<Origin, LogBuffer>,
/// Routes with their associated limiter sets.
routes: Vec<(Route, Mutex<LimiterSet>)>,
/// Overall rate limiters applied after route checks pass.
@@ -121,7 +121,7 @@ impl LogHandler {
let mut text = String::new();
let now = Utc::now().timestamp();
for (origin, buffer_mutex) in buffers.iter() {
for (origin, buffer) in buffers.iter() {
// Apply filter if present
if let Some(f) = filter {
if !origin.matches_filter(f) {
@@ -129,21 +129,13 @@ impl LogHandler {
}
}
// We can't await inside read_all, so use try_lock
// If the buffer is locked, skip it (rare case)
let Some(buffer) = buffer_mutex.try_lock().ok() else {
continue;
};
use std::fmt::Write;
writeln!(&mut text, "=== [{origin}] ===").unwrap();
writeln!(&mut text, "{} log messages (newest first):", buffer.len()).unwrap();
// Collect and reverse to show newest first
let messages: Vec<_> = buffer.iter().collect();
for log_msg in messages.into_iter().rev() {
buffer.for_each(|log_msg| {
self.write_log_msg(&mut text, log_msg, now);
}
});
text.push('\n');
}
@@ -176,26 +168,21 @@ impl LogHandler {
.log_buffers
.get_or_insert_with(
origin.clone(),
|| Mutex::new(CircularBuffer::new(buffer_size)),
|buffer_mutex| {
// Lock the buffer and process the message
let mut buffer = buffer_mutex.blocking_lock();
buffer.push_back(log_msg);
|| LogBuffer::new(buffer_size),
|buffer| {
if rate_limit_result.is_err() {
return None;
buffer.push_back(log_msg);
None
} else {
let mut text = String::default();
let now = Utc::now().timestamp();
buffer.push_back_and_drain(log_msg, |log_msg| {
self.write_log_msg(&mut text, log_msg, now);
});
Some(text)
}
let mut text = String::default();
let now = Utc::now().timestamp();
// Iterate in reverse (newest first) without copying
for log_msg in buffer.iter().rev() {
self.write_log_msg(&mut text, log_msg, now);
}
buffer.clear();
Some(text)
},
)
.await;
@@ -300,10 +287,7 @@ impl LogHandler {
}
// Check if message matches route's filter (if any)
let filter_matches = route
.filter
.as_ref()
.map_or(true, |f| f.matches(log_msg));
let filter_matches = route.filter.as_ref().map_or(true, |f| f.matches(log_msg));
if !filter_matches {
continue;
+2 -4
View File
@@ -34,6 +34,7 @@ use tokio_util::sync::CancellationToken;
use tracing::{debug, error, info, warn};
mod circular_buffer;
mod log_buffer;
mod log_handler;
use log_handler::{LogHandler, LogHandlerConfig};
@@ -563,10 +564,7 @@ impl Gateway {
async fn handle_gateway_command(&self, cmd: GatewayCommand) -> MessageHandlerResult {
match cmd {
GatewayCommand::Log { filter } => {
let text = self
.log_handler
.format_logs(filter.as_deref())
.await;
let text = self.log_handler.format_logs(filter.as_deref()).await;
Ok(AdminMessageResponse::new(text))
}
GatewayCommand::Query { query } => {