use appcore_core::redact_text_with_limit;
use parking_lot::{Mutex, RwLock};
use serde::{Deserialize, Serialize};
use std::collections::{BTreeMap, VecDeque};
use std::sync::Arc;
pub const MAX_OBSERVATION_ATTRIBUTES: usize = 32;
pub const MAX_OBSERVATION_NAME_BYTES: usize = 128;
pub const MAX_OBSERVATION_KEY_BYTES: usize = 64;
pub const MAX_OBSERVATION_VALUE_BYTES: usize = 1_024;
pub const MAX_OBSERVATION_TRACE_BYTES: usize = 256;
pub const MAX_IN_MEMORY_OBSERVATION_ITEMS: usize = 65_536;
pub const MAX_IN_MEMORY_OBSERVATION_BYTES: usize = 16 * 1024 * 1024;
pub const MAX_OBSERVATION_DRAINS: usize = 32;
const OBSERVATION_FIXED_BYTES: usize = std::mem::size_of::<ObservationEvent>()
+ std::mem::size_of::<Arc<ObservationEvent>>()
+ std::mem::size_of::<usize>() * 2;
const ATTRIBUTE_FIXED_BYTES: usize =
std::mem::size_of::<(String, String)>() + std::mem::size_of::<usize>() * 4;
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum ObservationKind {
Lifecycle,
Configuration,
Health,
Security,
Storage,
ControlPlane,
PeerRpc,
Scheduler,
Sync,
Audit,
Diagnostic,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum ObservationSeverity {
Debug,
Info,
Warning,
Error,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct ObservationEvent {
pub kind: ObservationKind,
pub severity: ObservationSeverity,
pub name: String,
pub timestamp_ms: u64,
pub trace_id: Option<String>,
pub attributes: BTreeMap<String, String>,
}
impl ObservationEvent {
pub fn new(
kind: ObservationKind,
severity: ObservationSeverity,
name: impl Into<String>,
timestamp_ms: u64,
) -> Self {
Self {
kind,
severity,
name: name.into(),
timestamp_ms,
trace_id: None,
attributes: BTreeMap::new(),
}
}
pub fn with_trace_id(mut self, trace_id: impl Into<String>) -> Self {
self.trace_id = Some(trace_id.into());
self
}
pub fn with_attribute(mut self, key: impl Into<String>, value: impl Into<String>) -> Self {
let key = key.into();
let sensitive = is_sensitive_key(&key);
let key = redact_text_with_limit(&key, MAX_OBSERVATION_KEY_BYTES);
if self.attributes.len() >= MAX_OBSERVATION_ATTRIBUTES
&& !self.attributes.contains_key(&key)
{
return self;
}
let value = if sensitive {
"[REDACTED]".to_string()
} else {
redact_text_with_limit(&value.into(), MAX_OBSERVATION_VALUE_BYTES)
};
self.attributes.insert(key, value);
self
}
pub(crate) fn redacted(mut self) -> Self {
self.name = redact_text_with_limit(&self.name, MAX_OBSERVATION_NAME_BYTES);
self.name.shrink_to_fit();
self.trace_id = self.trace_id.map(|value| {
let mut value = redact_text_with_limit(&value, MAX_OBSERVATION_TRACE_BYTES);
value.shrink_to_fit();
value
});
let mut attributes = BTreeMap::new();
for (key, value) in std::mem::take(&mut self.attributes)
.into_iter()
.take(MAX_OBSERVATION_ATTRIBUTES)
{
let sensitive = is_sensitive_key(&key);
let mut key = redact_text_with_limit(&key, MAX_OBSERVATION_KEY_BYTES);
let mut value = if sensitive {
"[REDACTED]".to_string()
} else {
redact_text_with_limit(&value, MAX_OBSERVATION_VALUE_BYTES)
};
key.shrink_to_fit();
value.shrink_to_fit();
attributes.insert(key, value);
}
self.attributes = attributes;
self
}
}
#[derive(Debug, Clone)]
pub struct SharedObservationEvent {
event: Arc<ObservationEvent>,
}
impl SharedObservationEvent {
pub fn new(event: ObservationEvent) -> Self {
Self {
event: Arc::new(event.redacted()),
}
}
pub fn as_event(&self) -> &ObservationEvent {
&self.event
}
fn clone_event_arc(&self) -> Arc<ObservationEvent> {
Arc::clone(&self.event)
}
}
pub trait ObservationSink: Send + Sync {
fn emit(&self, event: ObservationEvent);
fn emit_shared(&self, event: &SharedObservationEvent) {
self.emit(event.as_event().clone());
}
}
#[derive(Debug, Clone, Default)]
pub struct ObservationSnapshot {
events: Arc<VecDeque<Arc<ObservationEvent>>>,
}
impl ObservationSnapshot {
pub fn len(&self) -> usize {
self.events.len()
}
pub fn is_empty(&self) -> bool {
self.events.is_empty()
}
pub fn iter(&self) -> impl DoubleEndedIterator<Item = &ObservationEvent> {
self.events.iter().map(AsRef::as_ref)
}
pub fn recent(&self, limit: usize) -> impl Iterator<Item = &ObservationEvent> {
self.events
.iter()
.skip(self.events.len().saturating_sub(limit))
.map(AsRef::as_ref)
}
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
pub struct InMemoryObservationPressure {
pub entries: usize,
pub max_entries: usize,
pub used_bytes: usize,
pub peak_bytes: usize,
pub max_bytes: usize,
pub evictions: u64,
pub oversized_rejections: u64,
pub drain_rejections: u64,
}
#[derive(Debug)]
struct ObservationState {
events: Arc<VecDeque<Arc<ObservationEvent>>>,
pressure: InMemoryObservationPressure,
}
type ObservationDrain = Arc<dyn ObservationSink>;
type ObservationDrainGeneration = Arc<Vec<ObservationDrain>>;
#[derive(Clone)]
pub struct InMemoryObservationSink {
capacity: usize,
state: Arc<Mutex<ObservationState>>,
drains: Arc<RwLock<ObservationDrainGeneration>>,
}
impl std::fmt::Debug for InMemoryObservationSink {
fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
formatter
.debug_struct("InMemoryObservationSink")
.field("capacity", &self.capacity)
.field("event_count", &self.len())
.field("pressure", &self.pressure())
.field("drain_count", &self.drains.read().len())
.finish()
}
}
impl InMemoryObservationSink {
pub fn new(capacity: usize) -> Self {
let capacity = capacity.clamp(1, MAX_IN_MEMORY_OBSERVATION_ITEMS);
let max_bytes = default_max_bytes(capacity);
Self::with_limits(capacity, max_bytes)
}
pub fn with_max_bytes(capacity: usize, max_bytes: usize) -> Self {
let capacity = capacity.clamp(1, MAX_IN_MEMORY_OBSERVATION_ITEMS);
Self::with_limits(capacity, max_bytes.clamp(1, default_max_bytes(capacity)))
}
fn with_limits(capacity: usize, max_bytes: usize) -> Self {
Self {
capacity,
state: Arc::new(Mutex::new(ObservationState {
events: Arc::new(VecDeque::new()),
pressure: InMemoryObservationPressure {
max_entries: capacity,
max_bytes,
..InMemoryObservationPressure::default()
},
})),
drains: Arc::new(RwLock::new(Arc::new(Vec::new()))),
}
}
pub fn add_drain(&self, drain: Arc<dyn ObservationSink>) {
let _ = self.try_add_drain(drain);
}
pub fn try_add_drain(&self, drain: Arc<dyn ObservationSink>) -> bool {
let accepted = {
let mut drains = self.drains.write();
if drains.len() >= MAX_OBSERVATION_DRAINS {
false
} else {
Arc::make_mut(&mut drains).push(drain);
true
}
};
if !accepted {
let mut state = self.state.lock();
state.pressure.drain_rejections = state.pressure.drain_rejections.saturating_add(1);
}
accepted
}
pub fn drain_count(&self) -> usize {
self.drains.read().len()
}
pub fn snapshot(&self) -> Vec<ObservationEvent> {
self.shared_snapshot().iter().cloned().collect()
}
pub fn shared_snapshot(&self) -> ObservationSnapshot {
ObservationSnapshot {
events: Arc::clone(&self.state.lock().events),
}
}
pub fn pressure(&self) -> InMemoryObservationPressure {
self.state.lock().pressure
}
pub fn len(&self) -> usize {
self.state.lock().pressure.entries
}
pub fn is_empty(&self) -> bool {
self.len() == 0
}
}
impl Default for InMemoryObservationSink {
fn default() -> Self {
Self::new(1_024)
}
}
impl ObservationSink for InMemoryObservationSink {
fn emit(&self, event: ObservationEvent) {
self.emit_shared(&SharedObservationEvent::new(event));
}
fn emit_shared(&self, event: &SharedObservationEvent) {
let retained_bytes = observation_retained_bytes(event.as_event());
let mut state = self.state.lock();
if retained_bytes > state.pressure.max_bytes {
state.pressure.oversized_rejections =
state.pressure.oversized_rejections.saturating_add(1);
} else {
while state.pressure.entries >= self.capacity
|| state.pressure.used_bytes.saturating_add(retained_bytes)
> state.pressure.max_bytes
{
let Some(removed) = Arc::make_mut(&mut state.events).pop_front() else {
break;
};
state.pressure.entries = state.pressure.entries.saturating_sub(1);
state.pressure.used_bytes = state
.pressure
.used_bytes
.saturating_sub(observation_retained_bytes(&removed));
state.pressure.evictions = state.pressure.evictions.saturating_add(1);
}
Arc::make_mut(&mut state.events).push_back(event.clone_event_arc());
state.pressure.entries = state.pressure.entries.saturating_add(1);
state.pressure.used_bytes = state.pressure.used_bytes.saturating_add(retained_bytes);
state.pressure.peak_bytes = state.pressure.peak_bytes.max(state.pressure.used_bytes);
}
drop(state);
let drains = Arc::clone(&self.drains.read());
for drain in drains.iter() {
drain.emit_shared(event);
}
}
}
fn default_max_bytes(capacity: usize) -> usize {
capacity
.saturating_mul(maximum_observation_retained_bytes())
.clamp(1, MAX_IN_MEMORY_OBSERVATION_BYTES)
}
fn maximum_observation_retained_bytes() -> usize {
OBSERVATION_FIXED_BYTES
.saturating_add(MAX_OBSERVATION_NAME_BYTES)
.saturating_add(MAX_OBSERVATION_TRACE_BYTES)
.saturating_add(
MAX_OBSERVATION_ATTRIBUTES.saturating_mul(
ATTRIBUTE_FIXED_BYTES
.saturating_add(MAX_OBSERVATION_KEY_BYTES)
.saturating_add(MAX_OBSERVATION_VALUE_BYTES),
),
)
}
fn observation_retained_bytes(event: &ObservationEvent) -> usize {
OBSERVATION_FIXED_BYTES
.saturating_add(event.name.capacity())
.saturating_add(event.trace_id.as_ref().map_or(0, String::capacity))
.saturating_add(event.attributes.iter().fold(0usize, |total, (key, value)| {
total
.saturating_add(ATTRIBUTE_FIXED_BYTES)
.saturating_add(key.capacity())
.saturating_add(value.capacity())
}))
}
fn is_sensitive_key(key: &str) -> bool {
["secret", "password", "token", "credential", "private_key"]
.iter()
.any(|fragment| {
key.as_bytes()
.windows(fragment.len())
.any(|candidate| candidate.eq_ignore_ascii_case(fragment.as_bytes()))
})
}
#[cfg(test)]
#[path = "observation_tests.rs"]
mod tests;