diff --git a/signal-gateway-bin/src/json/mod.rs b/signal-gateway-bin/src/json/mod.rs index 0a543fa..d5675a3 100644 --- a/signal-gateway-bin/src/json/mod.rs +++ b/signal-gateway-bin/src/json/mod.rs @@ -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, /// Timestamp - accepts Unix epoch seconds (int or string) or RFC3339 string diff --git a/signal-gateway/src/concurrent_map.rs b/signal-gateway/src/concurrent_map.rs index 62c3f27..3764f73 100644 --- a/signal-gateway/src/concurrent_map.rs +++ b/signal-gateway/src/concurrent_map.rs @@ -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); diff --git a/signal-gateway/src/gateway/log_buffer.rs b/signal-gateway/src/gateway/log_buffer.rs new file mode 100644 index 0000000..4c5b527 --- /dev/null +++ b/signal-gateway/src/gateway/log_buffer.rs @@ -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>, +} + +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() + } +} diff --git a/signal-gateway/src/gateway/log_handler.rs b/signal-gateway/src/gateway/log_handler.rs index 43e21cb..03c8bd5 100644 --- a/signal-gateway/src/gateway/log_handler.rs +++ b/signal-gateway/src/gateway/log_handler.rs @@ -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, /// Log buffers keyed by origin (app + host). Lazily created. - log_buffers: ConcurrentMap>>, + log_buffers: ConcurrentMap, /// Routes with their associated limiter sets. routes: Vec<(Route, Mutex)>, /// 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; diff --git a/signal-gateway/src/gateway/mod.rs b/signal-gateway/src/gateway/mod.rs index cd39ef8..15576ff 100644 --- a/signal-gateway/src/gateway/mod.rs +++ b/signal-gateway/src/gateway/mod.rs @@ -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 } => {