add CollectedAt to message, get rid of backfilling timestamps with Utc::now()
This commit is contained in:
@@ -157,10 +157,8 @@ impl LogHandler {
|
|||||||
}
|
}
|
||||||
|
|
||||||
/// Consume a new log message from the given origin
|
/// Consume a new log message from the given origin
|
||||||
pub async fn handle_log_message(&self, mut log_msg: LogMessage, origin: Origin) {
|
pub async fn handle_log_message(&self, log_msg: LogMessage, origin: Origin) {
|
||||||
let ts_sec = *log_msg
|
let ts_sec = log_msg.get_timestamp_or_fallback();
|
||||||
.timestamp
|
|
||||||
.get_or_insert_with(|| Utc::now().timestamp());
|
|
||||||
|
|
||||||
let rate_limit_result = self.check_rate_limiters(&log_msg, &origin, ts_sec).await;
|
let rate_limit_result = self.check_rate_limiters(&log_msg, &origin, ts_sec).await;
|
||||||
|
|
||||||
|
|||||||
@@ -1,5 +1,6 @@
|
|||||||
//! Log message schema and types.
|
//! Log message schema and types.
|
||||||
|
|
||||||
|
use chrono::{DateTime, Utc};
|
||||||
use serde::{Deserialize, de};
|
use serde::{Deserialize, de};
|
||||||
|
|
||||||
/// Log severity level, following syslog conventions.
|
/// Log severity level, following syslog conventions.
|
||||||
@@ -94,10 +95,12 @@ impl<'de> Deserialize<'de> for Level {
|
|||||||
pub struct LogMessage {
|
pub struct LogMessage {
|
||||||
/// Severity level of the message.
|
/// Severity level of the message.
|
||||||
pub level: Level,
|
pub level: Level,
|
||||||
/// Unix timestamp in seconds.
|
/// Unix timestamp in seconds (from the log source).
|
||||||
pub timestamp: Option<i64>,
|
pub timestamp: Option<i64>,
|
||||||
/// Nanosecond component of the timestamp.
|
/// Nanosecond component of the timestamp.
|
||||||
pub timestamp_nanos: u32,
|
pub timestamp_nanos: u32,
|
||||||
|
/// When this message was collected by the gateway.
|
||||||
|
pub collected_at: DateTime<Utc>,
|
||||||
/// Hostname where the log originated.
|
/// Hostname where the log originated.
|
||||||
pub hostname: Option<Box<str>>,
|
pub hostname: Option<Box<str>>,
|
||||||
/// Application name that generated the log.
|
/// Application name that generated the log.
|
||||||
@@ -127,6 +130,13 @@ impl LogMessage {
|
|||||||
line: None,
|
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.
|
/// Builder for constructing [`LogMessage`] instances.
|
||||||
@@ -192,6 +202,7 @@ impl LogMessageBuilder {
|
|||||||
level: self.level,
|
level: self.level,
|
||||||
timestamp: self.timestamp,
|
timestamp: self.timestamp,
|
||||||
timestamp_nanos: self.timestamp_nanos,
|
timestamp_nanos: self.timestamp_nanos,
|
||||||
|
collected_at: Utc::now(),
|
||||||
hostname: self.hostname,
|
hostname: self.hostname,
|
||||||
appname: self.appname,
|
appname: self.appname,
|
||||||
msg: self.msg,
|
msg: self.msg,
|
||||||
|
|||||||
Reference in New Issue
Block a user