From 77a0767db56f9f463e8851ddd4273a931fc56c3b Mon Sep 17 00:00:00 2001 From: Chris Beck Date: Sat, 6 Dec 2025 19:30:01 -0700 Subject: [PATCH] add CollectedAt to message, get rid of backfilling timestamps with Utc::now() --- signal-gateway/src/gateway/log_handler.rs | 6 ++---- signal-gateway/src/log_message.rs | 13 ++++++++++++- 2 files changed, 14 insertions(+), 5 deletions(-) diff --git a/signal-gateway/src/gateway/log_handler.rs b/signal-gateway/src/gateway/log_handler.rs index 08e7eab..4722689 100644 --- a/signal-gateway/src/gateway/log_handler.rs +++ b/signal-gateway/src/gateway/log_handler.rs @@ -157,10 +157,8 @@ impl LogHandler { } /// Consume a new log message from the given origin - pub async fn handle_log_message(&self, mut log_msg: LogMessage, origin: Origin) { - let ts_sec = *log_msg - .timestamp - .get_or_insert_with(|| Utc::now().timestamp()); + pub async fn handle_log_message(&self, log_msg: LogMessage, origin: Origin) { + let ts_sec = log_msg.get_timestamp_or_fallback(); let rate_limit_result = self.check_rate_limiters(&log_msg, &origin, ts_sec).await; diff --git a/signal-gateway/src/log_message.rs b/signal-gateway/src/log_message.rs index 68555c4..b03ca12 100644 --- a/signal-gateway/src/log_message.rs +++ b/signal-gateway/src/log_message.rs @@ -1,5 +1,6 @@ //! Log message schema and types. +use chrono::{DateTime, Utc}; use serde::{Deserialize, de}; /// Log severity level, following syslog conventions. @@ -94,10 +95,12 @@ impl<'de> Deserialize<'de> for Level { pub struct LogMessage { /// Severity level of the message. pub level: Level, - /// Unix timestamp in seconds. + /// Unix timestamp in seconds (from the log source). pub timestamp: Option, /// Nanosecond component of the timestamp. pub timestamp_nanos: u32, + /// When this message was collected by the gateway. + pub collected_at: DateTime, /// Hostname where the log originated. pub hostname: Option>, /// Application name that generated the log. @@ -127,6 +130,13 @@ impl LogMessage { line: None, } } + + /// Get the timestamp in seconds, using the source timestamp if available, + /// otherwise falling back to the collection time. + pub fn get_timestamp_or_fallback(&self) -> i64 { + self.timestamp + .unwrap_or_else(|| self.collected_at.timestamp()) + } } /// Builder for constructing [`LogMessage`] instances. @@ -192,6 +202,7 @@ impl LogMessageBuilder { level: self.level, timestamp: self.timestamp, timestamp_nanos: self.timestamp_nanos, + collected_at: Utc::now(), hostname: self.hostname, appname: self.appname, msg: self.msg,