use chrono::Duration;
use serde_json::Value;
use std::sync::Mutex;
use crate::sparkv::{BTreeSparKV, Config as ConfigSparKV};
use super::LogLevel;
use super::err_log_entry::ErrorLogEntry;
use super::interface::{LogStorage, LogWriter, Loggable, composite_key};
use crate::app_types::{ApplicationName, PdpID};
use crate::bootstrap_config::log_config::MemoryLogConfig;
use crate::log::BaseLogEntry;
use crate::log::loggable_fn::LoggableFn;
mod memory_calc;
use memory_calc::calculate_memory_usage;
const STORAGE_MUTEX_EXPECT_MESSAGE: &str = "MemoryLogger storage mutex should unlock";
pub(crate) struct MemoryLogger {
storage: Mutex<BTreeSparKV<serde_json::Value>>,
log_level: LogLevel,
pdp_id: PdpID,
app_name: Option<ApplicationName>,
}
impl MemoryLogger {
pub(crate) fn new(
config: MemoryLogConfig,
log_level: LogLevel,
pdp_id: PdpID,
app_name: Option<ApplicationName>,
) -> Self {
let default_config: ConfigSparKV = ConfigSparKV::default();
let sparkv_config = ConfigSparKV {
default_ttl: Duration::new(
config.log_ttl.try_into().expect("u64 that fits in a i64"),
0,
)
.expect("a valid duration"),
max_items: config.max_items.unwrap_or(default_config.max_items),
max_item_size: config.max_item_size.unwrap_or(default_config.max_item_size),
earliest_expiration_eviction: true,
..Default::default()
};
MemoryLogger {
storage: Mutex::new(BTreeSparKV::with_config_and_sizer(
sparkv_config,
Some(calculate_memory_usage),
)),
log_level,
pdp_id,
app_name,
}
}
fn log_entry<T: Loggable>(&self, entry: &T) {
let entry_id = entry.get_id().to_string();
let index_keys = entry.get_index_keys();
let json = to_json_value(entry);
let err = {
let mut storage = self.storage.lock().expect(STORAGE_MUTEX_EXPECT_MESSAGE);
match storage.set(&entry_id, json, index_keys.as_slice()) {
Ok(()) => return,
Err(err) => err,
}
};
let err_entry = ErrorLogEntry::from_loggable(
entry,
format!("could not store LogEntry to memory: {err:?}"),
);
fallback::log(err_entry, &self.pdp_id, self.app_name.as_ref());
}
}
mod fallback {
use super::ErrorLogEntry;
use crate::LogLevel;
use crate::app_types::{ApplicationName, PdpID};
use crate::log::StdOutLoggerMode;
use crate::log::log_strategy::LogStrategyLogger;
use crate::log::stdout_logger::StdOutLogger;
pub(super) fn log(entry: ErrorLogEntry, pdp_id: &PdpID, app_name: Option<&ApplicationName>) {
use crate::log::interface::LogWriter;
let logger = StdOutLogger::new(LogLevel::TRACE, StdOutLoggerMode::Immediate);
let log_strategy = crate::log::LogStrategy::new_with_logger(
LogStrategyLogger::StdOut(logger),
*pdp_id,
app_name.cloned(),
None,
);
log_strategy.log_any(entry);
}
}
fn to_json_value<T: Loggable>(entry: &T) -> Value {
match serde_json::to_value(entry) {
Ok(json) => json,
Err(err) => {
let err_msg = format!("failed to serialize log entry to JSON: {err}");
serde_json::to_value(ErrorLogEntry::from_loggable(entry, err_msg.clone()))
.expect(&err_msg)
},
}
}
impl LogWriter for MemoryLogger {
fn log_any<T: Loggable>(&self, entry: T) {
if !entry.can_log(self.log_level) {
return;
}
self.log_entry(&entry);
}
fn log_fn<F, R>(&self, log_fn: LoggableFn<F>)
where
R: Loggable,
F: Fn(BaseLogEntry) -> R,
{
if log_fn.can_log(self.log_level) {
let entry = log_fn.build();
self.log_entry(&entry);
}
}
}
impl LogStorage for MemoryLogger {
fn pop_logs(&self) -> Vec<serde_json::Value> {
self.storage
.lock()
.expect(STORAGE_MUTEX_EXPECT_MESSAGE)
.drain()
.map(|(_k, value)| value)
.collect()
}
fn get_log_by_id(&self, id: &str) -> Option<serde_json::Value> {
self.storage
.lock()
.expect(STORAGE_MUTEX_EXPECT_MESSAGE)
.get(id)
.cloned()
}
fn get_log_ids(&self) -> Vec<String> {
self.storage
.lock()
.expect(STORAGE_MUTEX_EXPECT_MESSAGE)
.get_keys()
}
fn get_logs_by_tag(&self, tag: &str) -> Vec<serde_json::Value> {
self.storage
.lock()
.expect(STORAGE_MUTEX_EXPECT_MESSAGE)
.get_by_index_key(tag)
.map(std::borrow::ToOwned::to_owned)
.collect()
}
fn get_logs_by_request_id(&self, request_id: &str) -> Vec<serde_json::Value> {
self.storage
.lock()
.expect(STORAGE_MUTEX_EXPECT_MESSAGE)
.get_by_index_key(request_id)
.map(std::borrow::ToOwned::to_owned)
.collect()
}
fn get_logs_by_request_id_and_tag(
&self,
request_id: &str,
tag: &str,
) -> Vec<serde_json::Value> {
let key = composite_key(request_id, tag);
self.storage
.lock()
.expect(STORAGE_MUTEX_EXPECT_MESSAGE)
.get_by_index_key(&key)
.map(std::borrow::ToOwned::to_owned)
.collect()
}
}
#[cfg(test)]
mod tests {
use super::super::interface::Indexed;
use super::super::{AuthorizationLogInfo, LogEntry, LogType};
use super::*;
use crate::log::gen_uuid7;
use serde_json::json;
use test_utils::assert_eq;
fn create_memory_logger(pdp_id: PdpID, app_name: Option<ApplicationName>) -> MemoryLogger {
let config = MemoryLogConfig {
log_ttl: 60,
max_items: None,
max_item_size: None,
};
MemoryLogger::new(config, LogLevel::TRACE, pdp_id, app_name)
}
#[test]
fn test_log_and_get_logs() {
let pdp_id = PdpID::new();
let app_name = None;
let logger = create_memory_logger(pdp_id, app_name.clone());
let entry1 = LogEntry::new(BaseLogEntry::new_decision_opt_request_id(None))
.set_message("some message".to_string())
.set_auth_info(AuthorizationLogInfo {
action: "test_action".to_string(),
resource: "test_resource".to_string(),
context: serde_json::json!({}),
authorize_info: Vec::default(),
authorized: true,
entities: serde_json::json!({}),
});
let entry2 = LogEntry::new(BaseLogEntry::new_system_opt_request_id(
LogLevel::INFO,
None,
));
assert!(
entry1.base.id < entry2.base.id,
"entry1.base.id should be lower than in entry2"
);
logger.log_any(entry1.clone());
logger.log_any(entry2.clone());
let entry1_json = json!(entry1.clone());
let entry2_json = json!(entry2.clone());
assert_eq!(logger.get_log_ids().len(), 2);
assert_eq!(
logger.get_log_by_id(&entry1.get_id().to_string()).unwrap(),
entry1_json,
"Failed to get log entry by id"
);
assert_eq!(
logger.get_log_by_id(&entry2.get_id().to_string()).unwrap(),
entry2_json,
"Failed to get log entry by id"
);
let logs = logger.pop_logs();
assert_eq!(logs.len(), 2);
assert_eq!(logs[0], entry1_json, "First log entry is incorrect");
assert_eq!(logs[1], entry2_json, "Second log entry is incorrect");
assert!(
logger.get_log_ids().is_empty(),
"Logs were not fully popped"
);
}
#[test]
fn test_pop_logs() {
let pdp_id = PdpID::new();
let app_name = None;
let logger = create_memory_logger(pdp_id, app_name.clone());
let entry1 = LogEntry::new(BaseLogEntry::new_decision_opt_request_id(None));
let entry2 = LogEntry::new(BaseLogEntry::new_metric_opt_request_id(None));
logger.log_any(entry1.clone());
logger.log_any(entry2.clone());
let entry1_json = json!(entry1.clone());
let entry2_json = json!(entry2.clone());
let logs = logger.pop_logs();
assert_eq!(logs.len(), 2);
assert_eq!(logs[0], entry1_json, "First log entry is incorrect");
assert_eq!(logs[1], entry2_json, "Second log entry is incorrect");
assert!(
logger.get_log_ids().is_empty(),
"Logs were not fully popped"
);
}
#[test]
fn test_log_index() {
let request_id = gen_uuid7();
let logger = MemoryLogger::new(
MemoryLogConfig {
log_ttl: 10,
max_item_size: None,
max_items: None,
},
LogLevel::DEBUG,
PdpID::new(),
None,
);
let entry_decision = LogEntry::new(BaseLogEntry::new_decision_opt_request_id(None));
logger.log_any(entry_decision);
let entry_system_info = LogEntry::new(BaseLogEntry::new_system_opt_request_id(
LogLevel::INFO,
Some(request_id),
));
logger.log_any(entry_system_info);
let entry_system_debug = LogEntry::new(BaseLogEntry::new_system_opt_request_id(
LogLevel::DEBUG,
Some(request_id),
));
logger.log_any(entry_system_debug);
let entry_metric = LogEntry::new(BaseLogEntry::new_metric_opt_request_id(None));
logger.log_any(entry_metric);
let entry_system_warn = LogEntry::new(BaseLogEntry::new_system_opt_request_id(
LogLevel::WARN,
None,
));
logger.log_any(entry_system_warn);
assert!(
logger
.get_logs_by_request_id(request_id.to_string().as_str())
.len()
== 2,
"2 log entries should be present for request id: {request_id}"
);
assert!(
logger
.get_logs_by_request_id_and_tag(
request_id.to_string().as_str(),
LogLevel::DEBUG.to_string().as_str()
)
.len()
== 1,
"1 log entries should be present for request id: {request_id} and debug level"
);
assert!(
logger
.get_logs_by_tag(LogType::System.to_string().as_str())
.len()
== 3,
"3 system log entries should be present"
);
assert!(
logger
.get_logs_by_tag(LogLevel::WARN.to_string().as_str())
.len()
== 1,
"1 system log entry should be present with WARN level"
);
}
#[test]
fn test_max_items_config() {
let default_config: ConfigSparKV = ConfigSparKV::default();
let logger = MemoryLogger::new(
MemoryLogConfig {
log_ttl: 10,
max_items: None,
max_item_size: None,
},
LogLevel::DEBUG,
PdpID::new(),
None,
);
assert_eq!(
logger.storage.lock().unwrap().config.max_items,
default_config.max_items
);
let logger = MemoryLogger::new(
MemoryLogConfig {
log_ttl: 10,
max_items: Some(0),
max_item_size: None,
},
LogLevel::DEBUG,
PdpID::new(),
None,
);
assert_eq!(logger.storage.lock().unwrap().config.max_items, 0);
let logger = MemoryLogger::new(
MemoryLogConfig {
log_ttl: 10,
max_items: Some(500),
max_item_size: None,
},
LogLevel::DEBUG,
PdpID::new(),
None,
);
assert_eq!(logger.storage.lock().unwrap().config.max_items, 500);
}
#[test]
fn test_max_item_size_config() {
let default_config: ConfigSparKV = ConfigSparKV::default();
let logger = MemoryLogger::new(
MemoryLogConfig {
log_ttl: 10,
max_items: None,
max_item_size: None,
},
LogLevel::DEBUG,
PdpID::new(),
None,
);
assert_eq!(
logger.storage.lock().unwrap().config.max_item_size,
default_config.max_item_size
);
let logger = MemoryLogger::new(
MemoryLogConfig {
log_ttl: 10,
max_items: None,
max_item_size: Some(0),
},
LogLevel::DEBUG,
PdpID::new(),
None,
);
assert_eq!(logger.storage.lock().unwrap().config.max_item_size, 0);
let logger = MemoryLogger::new(
MemoryLogConfig {
log_ttl: 10,
max_items: None,
max_item_size: Some(10_000),
},
LogLevel::DEBUG,
PdpID::new(),
None,
);
assert_eq!(logger.storage.lock().unwrap().config.max_item_size, 10_000);
}
#[test]
fn test_capacity_overflow_evicts_oldest_and_keeps_logging() {
let max_items = 3;
let logger = MemoryLogger::new(
MemoryLogConfig {
log_ttl: 60,
max_items: Some(max_items),
max_item_size: None,
},
LogLevel::TRACE,
PdpID::new(),
None,
);
let total = max_items * 4;
let mut entries = Vec::with_capacity(total);
for _ in 0..total {
let entry = LogEntry::new(BaseLogEntry::new_decision_opt_request_id(None));
logger.log_any(entry.clone());
entries.push(entry);
}
let ids = logger.get_log_ids();
assert_eq!(
ids.len(),
max_items,
"storage should cap at max_items={max_items}, got {}",
ids.len()
);
for entry in entries.iter().rev().take(max_items) {
let id = entry.get_id().to_string();
assert!(
logger.get_log_by_id(&id).is_some(),
"expected most-recent entry {id} to be retained after overflow"
);
}
}
}