From 498e7c7c1e08432c0c16ab602304d0f03dc11b4c Mon Sep 17 00:00:00 2001 From: Chris Beck Date: Tue, 16 Dec 2025 02:06:40 -0700 Subject: [PATCH] refactor path prefix handling, normalize log paths at ingestion time not at formatting time --- signal-gateway-bin/src/main.rs | 78 +++++++++++--------- signal-gateway-log-ingest/src/lib.rs | 1 + signal-gateway-log-ingest/src/path_prefix.rs | 37 +++++++++- signal-gateway/src/gateway/log_handler.rs | 3 +- signal-gateway/src/gateway/mod.rs | 29 +++++++- signal-gateway/src/log_format.rs | 19 +---- signal-gateway/src/rate_limiter.rs | 2 +- 7 files changed, 111 insertions(+), 58 deletions(-) diff --git a/signal-gateway-bin/src/main.rs b/signal-gateway-bin/src/main.rs index 8088b62..b6cc178 100644 --- a/signal-gateway-bin/src/main.rs +++ b/signal-gateway-bin/src/main.rs @@ -21,7 +21,7 @@ use app_code::AppCodeConfigExt; mod listen_http; use listen_http::start_http_task; -use signal_gateway_log_ingest::{JsonConfig, SyslogConfig}; +use signal_gateway_log_ingest::{JsonConfig, StripPathPrefixes, SyslogConfig}; /// Top-level configuration for signal-gateway. #[derive(Conf, Debug)] @@ -53,6 +53,8 @@ pub struct Config { claude: Option, #[conf(flatten, serde(flatten))] gateway: GatewayConfig, + #[conf(long, env, value_parser = serde_json::from_str, default_value = "[\"/home/*/\", \".cargo/registry/src/*/\"]")] + remove_path_prefixes: Vec, } fn init_logging() { @@ -111,26 +113,27 @@ async fn main() -> Result<(), Box> { let token = CancellationToken::new(); // Build the command router - let mut router_builder = CommandRouter::builder() - .route("--help", Handling::Help) - .route("-h", Handling::Help); + let command_router = { + let mut router_builder = CommandRouter::builder() + .route("--help", Handling::Help) + .route("-h", Handling::Help); - // Add admin HTTP handler with its configured prefix - if let Some(admin_http) = config.admin_http { - let prefix = admin_http.command_prefix.clone(); - let handler = admin_http.into_handler(); - router_builder = router_builder.route(prefix, Handling::Custom(handler)); - } + // Add admin HTTP handler with its configured prefix + if let Some(admin_http) = config.admin_http { + let prefix = admin_http.command_prefix.clone(); + let handler = admin_http.into_handler(); + router_builder = router_builder.route(prefix, Handling::Custom(handler)); + } - // Add gateway commands for "/" prefix - router_builder = router_builder.route("/", Handling::GatewayCommand); + // Add gateway commands for "/" prefix + router_builder = router_builder.route("/", Handling::GatewayCommand); - // Add Claude as default handler if configured - if config.claude.is_some() { - router_builder = router_builder.route("", Handling::Claude); - } - - let command_router = router_builder.build(); + // Add Claude as default handler if configured + if config.claude.is_some() { + router_builder = router_builder.route("", Handling::Claude); + } + router_builder.build() + }; // Build AppCode tools if configured let app_code_tools = if !config.app_code.is_empty() { @@ -150,25 +153,32 @@ async fn main() -> Result<(), Box> { None }; - let mut gateway_builder = Gateway::builder(config.gateway) - .with_cancellation_token(token.clone()) - .with_command_router(command_router); + // Build path prefix stripper + let path_stripper = + StripPathPrefixes::new(config.remove_path_prefixes.iter().map(String::as_str)); - if let Some(tools) = app_code_tools { - gateway_builder = gateway_builder.with_tools(tools); - } + let gateway = { + let mut gateway_builder = Gateway::builder(config.gateway) + .with_cancellation_token(token.clone()) + .with_command_router(command_router) + .with_path_normalization_fn(move |p| path_stripper.strip_path_prefix(p)); - // Add Claude assistant if configured - if let Some(claude_config) = config.claude { - gateway_builder = gateway_builder.with_assistant(move |tool_executor| { - Box::new( - ClaudeAssistant::new(claude_config, tool_executor) - .expect("Failed to initialize Claude assistant"), - ) - }); - } + if let Some(tools) = app_code_tools { + gateway_builder = gateway_builder.with_tools(tools); + } - let gateway = gateway_builder.build().await; + // Add Claude assistant if configured + if let Some(claude_config) = config.claude { + gateway_builder = gateway_builder.with_assistant(move |tool_executor| { + Box::new( + ClaudeAssistant::new(claude_config, tool_executor) + .expect("Failed to initialize Claude assistant"), + ) + }); + } + + gateway_builder.build().await + }; let listener = TcpListener::bind(config.http_listen_addr).await.unwrap(); info!("Listening for http on {}", config.http_listen_addr); diff --git a/signal-gateway-log-ingest/src/lib.rs b/signal-gateway-log-ingest/src/lib.rs index 50468a3..2175b93 100644 --- a/signal-gateway-log-ingest/src/lib.rs +++ b/signal-gateway-log-ingest/src/lib.rs @@ -9,4 +9,5 @@ pub mod path_prefix; pub mod syslog; pub use json::JsonConfig; +pub use path_prefix::StripPathPrefixes; pub use syslog::SyslogConfig; diff --git a/signal-gateway-log-ingest/src/path_prefix.rs b/signal-gateway-log-ingest/src/path_prefix.rs index e4af982..0bff22f 100644 --- a/signal-gateway-log-ingest/src/path_prefix.rs +++ b/signal-gateway-log-ingest/src/path_prefix.rs @@ -1,7 +1,39 @@ use globset::{Glob, GlobMatcher}; +pub struct StripPathPrefixes { + finders: Vec, +} + +impl StripPathPrefixes { + pub fn new<'a, T>(iter: T) -> Self + where + T: Iterator, + { + Self { + finders: iter.map(Finder::new).collect(), + } + } + + pub fn strip_path_prefix<'a>(&self, mut target: &'a str) -> &'a str { + for f in self.finders.iter() { + target = f.strip_path_prefix(target); + } + target + } +} + +impl> FromIterator for StripPathPrefixes { + fn from_iter(iter: T) -> Self + where + T: IntoIterator, + { + let finders = iter.into_iter().map(|s| Finder::new(s.as_ref())).collect(); + Self { finders } + } +} + /// Trait for object which can strip prefixes from a path -pub trait PathPrefixFinder { +trait PathPrefixFinder { /// Find the largest index such that target[..idx] matches to glob fn find_path_prefix(&self, target: &str) -> Option; /// Strip a glob pattern from a target path @@ -15,6 +47,7 @@ pub trait PathPrefixFinder { } /// Matches a glob pattern to a path and strips matching prefix +#[derive(Clone)] pub struct Finder { inner: FinderInner, } @@ -34,6 +67,7 @@ impl PathPrefixFinder for Finder { } } +#[derive(Clone)] enum FinderInner { LeadingDoubleStar { rneedle: Box, @@ -50,6 +84,7 @@ enum FinderInner { }, } +#[derive(Clone)] struct GlobSegment { prefix: Box, matcher: GlobMatcher, diff --git a/signal-gateway/src/gateway/log_handler.rs b/signal-gateway/src/gateway/log_handler.rs index d15bed4..bed210b 100644 --- a/signal-gateway/src/gateway/log_handler.rs +++ b/signal-gateway/src/gateway/log_handler.rs @@ -167,7 +167,8 @@ impl LogHandler { } /// Consume a new log message from the given origin - pub async fn handle_log_message(&self, log_msg: LogMessage, origin: Origin) { + pub async fn handle_log_message(&self, log_msg: LogMessage) { + let origin = Origin::from(&log_msg); let rate_limit_result = self.check_rate_limiters(&log_msg, &origin).await; if let Err(reason) = &rate_limit_result { diff --git a/signal-gateway/src/gateway/mod.rs b/signal-gateway/src/gateway/mod.rs index 07e1784..932933e 100644 --- a/signal-gateway/src/gateway/mod.rs +++ b/signal-gateway/src/gateway/mod.rs @@ -253,6 +253,8 @@ pub struct Gateway { assistant: OnceLock, /// Additional tool executors added via the builder. extra_tool_executors: Vec>, + /// Path normalization function + path_normalization_fn: Option &str + Send + Sync>>, } /// Type alias for the assistant factory function. @@ -266,6 +268,7 @@ pub struct GatewayBuilder { command_router: Option, extra_tool_executors: Vec>, assistant_factory: Option, + path_normalization_fn: Option &str + Send + Sync>>, } impl GatewayBuilder { @@ -277,6 +280,7 @@ impl GatewayBuilder { command_router: None, extra_tool_executors: Vec::new(), assistant_factory: None, + path_normalization_fn: None, } } @@ -313,6 +317,15 @@ impl GatewayBuilder { self } + /// Set a path normalization function for the gateway, for source location of logs + pub fn with_path_normalization_fn(mut self, normalize_fn: F) -> Self + where + F: Fn(&str) -> &str + Send + Sync + 'static, + { + self.path_normalization_fn = Some(Box::new(normalize_fn)); + self + } + /// Build the gateway. pub async fn build(self) -> Arc { Gateway::new_internal( @@ -321,6 +334,7 @@ impl GatewayBuilder { self.command_router.unwrap_or_default(), self.extra_tool_executors, self.assistant_factory, + self.path_normalization_fn, ) .await } @@ -339,6 +353,7 @@ impl Gateway { command_router: CommandRouter, extra_tool_executors: Vec>, assistant_factory: Option, + path_normalization_fn: Option &str + Send + Sync>>, ) -> Arc { let child_token = token.child_token(); let (signal_alert_mq_tx, signal_alert_mq_rx) = unbounded_channel(); @@ -362,6 +377,7 @@ impl Gateway { command_router, assistant: OnceLock::new(), extra_tool_executors, + path_normalization_fn, }); // Initialize the assistant agent using the factory if one was provided @@ -912,9 +928,16 @@ impl Gateway { /// Process an incoming log message, buffering it and potentially triggering an alert. pub async fn handle_log_message(&self, log_msg: impl Into) { - let log_msg = log_msg.into(); - let origin = Origin::from(&log_msg); - self.log_handler.handle_log_message(log_msg, origin).await; + let mut log_msg = log_msg.into(); + if let Some(n_fn) = self.path_normalization_fn.as_ref() + && let Some(file) = log_msg.file.as_ref() + { + let normalized = (n_fn)(file); + if normalized.len() < file.len() { + log_msg.file = Some(normalized.to_owned().into_boxed_str()); + } + } + self.log_handler.handle_log_message(log_msg).await; } } diff --git a/signal-gateway/src/log_format.rs b/signal-gateway/src/log_format.rs index b404aa0..eb2e554 100644 --- a/signal-gateway/src/log_format.rs +++ b/signal-gateway/src/log_format.rs @@ -78,11 +78,7 @@ impl LogFormatConfig { if self.format_source_location { if let Some(file) = log_msg.file.as_ref() { - // Strip /home/{username}/ prefix if present - let trimmed = strip_prefix_and_one_slash(file, "/home/"); - // Strip .cargo/registry/src/{hash}/ if present - let trimmed = strip_prefix_and_one_slash(trimmed, ".cargo/registry/src/"); - write!(writer, "{trimmed}:")?; + write!(writer, "{file}:")?; if let Some(line) = log_msg.line.as_ref() { write!(writer, "{line}")?; } else { @@ -100,19 +96,6 @@ impl LogFormatConfig { } } -// Strip a prefix, then find the first remaining slash and skip up to that as well. -fn strip_prefix_and_one_slash<'a>(target: &'a str, prefix: &str) -> &'a str { - let Some(target) = target.strip_prefix(prefix) else { - return target; - }; - - if let Some((_, after)) = target.split_once('/') { - after - } else { - target - } -} - /// Write a relative timestamp like "T-30s", "T-1m30s", "T+5s". /// /// For durations < 1 second, shows milliseconds. diff --git a/signal-gateway/src/rate_limiter.rs b/signal-gateway/src/rate_limiter.rs index 874cf3a..8ef2cad 100644 --- a/signal-gateway/src/rate_limiter.rs +++ b/signal-gateway/src/rate_limiter.rs @@ -713,7 +713,7 @@ mod tests { #[test] fn source_location_rate_limiter_cleanup() { - // Cleanup triggers when len >= 2 * last_cleanup_size (starts at 8, so triggers at 16) + // Cleanup triggers when len >= 2 * last_cleanup_size (starts at 0, triggers at 2, 4, 8, 16...) let threshold = RateThreshold { times: 1, duration: Duration::from_secs(600),