use std::collections::BTreeMap;
use serde_json::Value;
use crate::config::TelemetryConfig;
use crate::fingerprint::compute_error_fingerprint;
use crate::harden::{clean_key, cleaned_key_claims_slot, harden_value, HardenLimits};
use crate::pii::{detect_secret_in_string, sanitize_payload, REDACTED_SENTINEL};
use crate::runtime::get_runtime_config;
use crate::schema::{event_name, get_strict_schema, validate_required_keys};
use super::LogEvent;
fn runtime_config_or_env() -> Option<TelemetryConfig> {
match get_runtime_config() {
Some(cfg) => Some(cfg),
None => TelemetryConfig::from_env().ok(),
}
}
fn first_context_string<'a>(event: &'a LogEvent, keys: &[&str]) -> Option<&'a str> {
for key in keys {
if let Some(Value::String(value)) = event.context.get(*key) {
return Some(value.as_str());
}
}
None
}
fn runtime_schema_error(event: &LogEvent) -> Option<String> {
let cfg = get_runtime_config()?;
match validate_required_keys(&event.context, &cfg.event_schema.required_keys) {
Ok(()) => None,
Err(err) => Some(err.message),
}
}
fn is_priority_key(key: &str) -> bool {
matches!(
key,
"service"
| "env"
| "version"
| "trace_id"
| "span_id"
| "session_id"
| "logger_name"
| "domain"
| "action"
| "resource"
| "status"
| "error_fingerprint"
)
}
pub(super) fn process_event(event: &mut LogEvent) {
let cfg = runtime_config_or_env();
let pii_max_depth = cfg.as_ref().map_or(8, |c| c.pii_max_depth);
let limits = HardenLimits {
max_value_length: cfg
.as_ref()
.map_or(1024, |c| c.security.max_attr_value_length),
max_attr_count: cfg.as_ref().map_or(64, |c| c.security.max_attr_count),
max_depth: cfg.as_ref().map_or(8, |c| c.security.max_nesting_depth),
};
extract_dars_fields(event);
inject_logger_name(event);
harden_input(event, limits);
add_error_fingerprint(event);
sanitize_context(event, pii_max_depth);
enforce_schema(event);
}
fn extract_dars_fields(event: &mut LogEvent) {
let meta = match event.event_metadata.as_ref() {
Some(meta) => meta,
None => return,
};
event
.context
.insert("domain".to_string(), Value::String(meta.domain.clone()));
event
.context
.insert("action".to_string(), Value::String(meta.action.clone()));
if let Some(resource) = meta.resource.as_ref() {
event
.context
.insert("resource".to_string(), Value::String(resource.clone()));
}
event
.context
.insert("status".to_string(), Value::String(meta.status.clone()));
}
fn inject_logger_name(event: &mut LogEvent) {
if event.target.is_empty() {
return;
}
if event.context.contains_key("logger_name") {
return;
}
event.context.insert(
"logger_name".to_string(),
Value::String(event.target.clone()),
);
}
fn harden_input(event: &mut LogEvent, limits: HardenLimits) {
for value in event.context.values_mut() {
*value = harden_value(value, limits, 1);
}
cap_attr_count(event, limits.max_attr_count);
harden_keys(&mut event.context);
}
fn harden_keys(context: &mut BTreeMap<String, Value>) {
use std::collections::BTreeSet;
let original = std::mem::take(context);
let mut verbatim: BTreeSet<String> = BTreeSet::new();
for (key, value) in original {
let name = clean_key(&key);
let untouched = name == key;
if !cleaned_key_claims_slot(
context.contains_key(&name),
untouched,
verbatim.contains(&name),
) {
continue;
}
context.insert(name.clone(), value);
if untouched {
verbatim.insert(name);
}
}
}
fn cap_attr_count(event: &mut LogEvent, max_attr_count: usize) {
if max_attr_count == 0 {
return;
}
if event.context.len() <= max_attr_count {
return;
}
use std::collections::BTreeSet;
let mut keep = BTreeSet::new();
let priority_keys: Vec<String> = event.context.keys().cloned().collect();
for key in &priority_keys {
if is_priority_key(key) {
keep.insert(key.clone());
}
}
let candidate_keys: Vec<String> = event.context.keys().cloned().collect();
for key in &candidate_keys {
if keep.len() >= max_attr_count {
break;
}
if is_priority_key(key) {
continue;
}
keep.insert(key.clone());
}
let original_context = std::mem::take(&mut event.context);
let mut retained = BTreeMap::new();
for (key, value) in original_context {
if keep.contains(&key) {
retained.insert(key, value);
}
}
event.context = retained;
}
fn add_error_fingerprint(event: &mut LogEvent) {
const ERROR_LEVELS: &[&str] = &["ERROR", "CRITICAL", "FATAL"];
if !ERROR_LEVELS.contains(&event.level.as_str()) {
return;
}
let error_name = first_context_string(event, &["error", "error_type", "exception"]);
let error_name = match error_name {
Some(error_name) => error_name,
None => return,
};
let stack = first_context_string(event, &["stack", "stacktrace"]);
let fingerprint = compute_error_fingerprint(error_name, stack);
event
.context
.insert("error_fingerprint".to_string(), Value::String(fingerprint));
}
fn sanitize_context(event: &mut LogEvent, max_depth: usize) {
let message_has_secret = detect_secret_in_string(&event.message);
if message_has_secret {
event.message = REDACTED_SENTINEL.to_string();
}
if event.context.is_empty() {
return;
}
let payload = Value::Object(event.context.clone().into_iter().collect());
let cleaned = sanitize_payload(&payload, true, max_depth);
let object = cleaned
.as_object()
.expect("sanitize_payload preserves object shape")
.clone();
event.context = object.into_iter().collect();
}
fn enforce_schema(event: &mut LogEvent) {
if let Some(message) = runtime_schema_error(event) {
event
.context
.insert("_schema_error".to_string(), Value::String(message));
return;
}
if !get_strict_schema() {
return;
}
let segments: Vec<&str> = event.message.split('.').collect();
match event_name(&segments) {
Ok(_) => {}
Err(_) => {
event.context.insert(
"_schema_error".to_string(),
Value::String(format!(
"event name {:?} does not match strict schema",
event.message
)),
);
}
}
}
#[cfg(test)]
#[path = "processors_tests.rs"]
mod tests;
#[cfg(test)]
#[path = "processors_message_pii_tests.rs"]
mod message_pii_tests;
#[cfg(test)]
#[path = "processors_edge_tests.rs"]
mod edge_tests;
#[cfg(test)]
#[path = "processors_key_hardening_tests.rs"]
mod key_hardening_tests;