use std::collections::HashMap;
use std::time::Duration;
use opentelemetry::InstrumentationScope;
use opentelemetry::logs::{LogRecord, Logger, LoggerProvider};
use opentelemetry_sdk::error::OTelSdkResult;
use opentelemetry_sdk::logs::{LogProcessor, SdkLogRecord, SdkLogger, SdkLoggerProvider};
pub(super) struct UniqueAttributesProcessor {
blank: SdkLogger,
}
impl UniqueAttributesProcessor {
pub(super) fn new() -> Self {
Self {
blank: SdkLoggerProvider::builder()
.build()
.logger("aviso-unique-attributes"),
}
}
fn rebuilt(&self, record: &SdkLogRecord) -> SdkLogRecord {
let attributes: Vec<_> = record.attributes_iter().collect();
let mut last = HashMap::with_capacity(attributes.len());
for (index, (key, _)) in attributes.iter().enumerate() {
last.insert(key.clone(), index);
}
let mut copy = self.blank.create_log_record();
if let Some(name) = record.event_name() {
copy.set_event_name(name);
}
if let Some(target) = record.target() {
copy.set_target(target.clone());
}
if let Some(timestamp) = record.timestamp() {
copy.set_timestamp(timestamp);
}
if let Some(observed) = record.observed_timestamp() {
copy.set_observed_timestamp(observed);
}
if let Some(context) = record.trace_context() {
copy.set_trace_context(context.trace_id, context.span_id, context.trace_flags);
}
if let Some(text) = record.severity_text() {
copy.set_severity_text(text);
}
if let Some(number) = record.severity_number() {
copy.set_severity_number(number);
}
if let Some(body) = record.body() {
copy.set_body(body.clone());
}
for (index, (key, value)) in attributes.into_iter().enumerate() {
if last.get(key) == Some(&index) {
copy.add_attribute(key.clone(), value.clone());
}
}
copy
}
}
fn has_repeated_key(record: &SdkLogRecord) -> bool {
let mut seen = std::collections::HashSet::new();
record.attributes_iter().any(|(key, _)| !seen.insert(key))
}
impl LogProcessor for UniqueAttributesProcessor {
fn emit(&self, record: &mut SdkLogRecord, _instrumentation: &InstrumentationScope) {
if has_repeated_key(record) {
*record = self.rebuilt(record);
}
}
fn force_flush(&self) -> OTelSdkResult {
Ok(())
}
fn shutdown_with_timeout(&self, _timeout: Duration) -> OTelSdkResult {
Ok(())
}
}
impl std::fmt::Debug for UniqueAttributesProcessor {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("UniqueAttributesProcessor")
.finish_non_exhaustive()
}
}
#[cfg(test)]
mod tests {
use super::*;
use opentelemetry::logs::{AnyValue, Severity};
use opentelemetry::trace::{SpanId, TraceFlags, TraceId};
use std::time::SystemTime;
fn record() -> SdkLogRecord {
SdkLoggerProvider::builder()
.build()
.logger("test")
.create_log_record()
}
fn emitted(record: &mut SdkLogRecord) {
UniqueAttributesProcessor::new()
.emit(record, &InstrumentationScope::builder("test").build());
}
fn attributes(record: &SdkLogRecord) -> Vec<(String, AnyValue)> {
record
.attributes_iter()
.map(|(key, value)| (key.to_string(), value.clone()))
.collect()
}
#[test]
fn record_without_repeated_keys_is_unchanged() {
let mut original = record();
original.add_attribute("request_id", "r1");
original.add_attribute("username", "alice");
let mut forwarded = original.clone();
emitted(&mut forwarded);
assert_eq!(attributes(&forwarded), attributes(&original));
}
#[test]
fn repeated_key_keeps_its_last_value_and_every_other_field() {
let timestamp = SystemTime::UNIX_EPOCH + Duration::from_secs(10);
let observed = SystemTime::UNIX_EPOCH + Duration::from_secs(11);
let trace_id = TraceId::from_hex("0af7651916cd43dd8448eb211c80319c").unwrap();
let span_id = SpanId::from_hex("b7ad6b7169203331").unwrap();
let mut rec = record();
rec.set_event_name("event src/routes/replay.rs:1");
rec.set_target("aviso_server::routes::replay");
rec.set_timestamp(timestamp);
rec.set_observed_timestamp(observed);
rec.set_trace_context(trace_id, span_id, Some(TraceFlags::SAMPLED));
rec.set_severity_text("INFO");
rec.set_severity_number(Severity::Info);
rec.set_body(AnyValue::String("replay started".into()));
rec.add_attribute("request_id", "r1");
rec.add_attribute("event_type", "outer");
rec.add_attribute("username", "alice");
rec.add_attribute("event_type", "inner");
rec.add_attribute("request_id", "r1");
rec.add_attribute("event_type", "event");
emitted(&mut rec);
assert_eq!(
attributes(&rec),
vec![
("username".to_string(), AnyValue::from("alice")),
("request_id".to_string(), AnyValue::from("r1")),
("event_type".to_string(), AnyValue::from("event")),
]
);
assert_eq!(rec.event_name(), Some("event src/routes/replay.rs:1"));
assert_eq!(
rec.target().map(AsRef::as_ref),
Some("aviso_server::routes::replay")
);
assert_eq!(rec.timestamp(), Some(timestamp));
assert_eq!(rec.observed_timestamp(), Some(observed));
let context = rec.trace_context().expect("trace context kept");
assert_eq!(context.trace_id, trace_id);
assert_eq!(context.span_id, span_id);
assert_eq!(context.trace_flags, Some(TraceFlags::SAMPLED));
assert_eq!(rec.severity_text(), Some("INFO"));
assert_eq!(rec.severity_number(), Some(Severity::Info));
assert_eq!(rec.body(), Some(&AnyValue::String("replay started".into())));
}
}