use crate::{
config::processed::OverflowPolicy,
error_handling::{InternalErrorReport, InternalErrorSource},
model::LogEvent,
subscriber::actor::{ActorAction, AppenderActor},
};
use fibre::{
error::TrySendError as FibreTrySendError, mpsc::BoundedSyncSender as FibreMpscBoundedSender,
};
use tracing_core::{metadata::LevelFilter, Metadata};
pub(crate) struct EventProcessor {
actors: Vec<AppenderActor>,
error_tx: Option<FibreMpscBoundedSender<InternalErrorReport>>,
max_level: LevelFilter,
}
impl EventProcessor {
pub(crate) fn new(
actors: Vec<AppenderActor>,
error_tx: Option<FibreMpscBoundedSender<InternalErrorReport>>,
) -> Self {
let max_level = actors
.iter()
.map(|a| a.filter.max_level())
.max()
.unwrap_or(LevelFilter::OFF);
Self {
actors,
error_tx,
max_level,
}
}
pub(crate) fn max_level(&self) -> LevelFilter {
self.max_level
}
pub(crate) fn event_enabled(&self, metadata: &Metadata<'_>) -> bool {
self.actors.iter().any(|a| a.filter.enabled(metadata))
}
pub(crate) fn close_channels(&self) {
for actor in &self.actors {
match &actor.action {
ActorAction::SendBytes(sender) => {
let _ = sender.close();
}
ActorAction::SendEvent(sender) => {
let _ = sender.close();
}
}
}
if let Some(tx) = &self.error_tx {
let _ = tx.close();
}
}
pub(crate) fn process_event(&self, event: LogEvent, metadata: &Metadata<'_>) {
let event_level = *metadata.level();
let rules: Vec<Option<(&str, &(LevelFilter, bool))>> = self
.actors
.iter()
.map(|actor| actor.filter.find_most_specific_rule(metadata))
.collect();
let mut winner: Option<(&str, bool)> = None;
for (prefix, (_, additive)) in rules.iter().flatten() {
if winner.map_or(true, |(wp, _)| prefix.len() > wp.len()) {
winner = Some((*prefix, *additive));
}
}
let non_additive_gate: Option<&str> = match winner {
Some((prefix, false)) => Some(prefix),
_ => None,
};
let mut event_senders: Vec<(&AppenderActor, &FibreMpscBoundedSender<LogEvent>)> = Vec::new();
for (actor, rule) in self.actors.iter().zip(&rules) {
if let Some(gate_prefix) = non_additive_gate {
let wired_to_gate = rule.map_or(false, |(prefix, _)| prefix == gate_prefix);
if !wired_to_gate {
continue;
}
}
let enabled = match rule {
Some((_, (level_filter, _))) => event_level <= *level_filter,
None => event_level <= actor.filter.default_level,
};
if !enabled {
continue;
}
match &actor.action {
ActorAction::SendBytes(sender) => match actor.formatter.format_event(&event) {
Ok(formatted_bytes) => self.send_bytes(actor, sender, formatted_bytes),
Err(e) => {
self.send_error_report(
InternalErrorSource::EventFormatting {
appender_name: actor.name.clone(),
},
e,
Some(format!("Event target: {}", event.target)),
);
}
},
ActorAction::SendEvent(sender) => event_senders.push((actor, sender)),
}
}
let last = event_senders.len().saturating_sub(1);
let mut event = Some(event);
for (i, (actor, sender)) in event_senders.into_iter().enumerate() {
let to_send = if i == last {
event.take().expect("event consumed before last sender")
} else {
event.as_ref().expect("event consumed early").clone()
};
self.send_event(actor, sender, to_send);
}
}
fn send_bytes(
&self,
actor: &AppenderActor,
sender: &FibreMpscBoundedSender<Vec<u8>>,
bytes: Vec<u8>,
) {
match actor.overflow {
OverflowPolicy::Block => {
let _ = sender.send(bytes);
}
OverflowPolicy::DropNewest => {
if let Err(FibreTrySendError::Full(_)) = sender.try_send(bytes) {
actor.drops.record(&actor.name);
}
}
}
}
fn send_event(
&self,
actor: &AppenderActor,
sender: &FibreMpscBoundedSender<LogEvent>,
event: LogEvent,
) {
match actor.overflow {
OverflowPolicy::Block => {
let _ = sender.send(event);
}
OverflowPolicy::DropNewest => {
if let Err(FibreTrySendError::Full(_)) = sender.try_send(event) {
actor.drops.record(&actor.name);
}
}
}
}
fn send_error_report<E: std::error::Error + 'static>(
&self,
source: InternalErrorSource,
error: E,
context: Option<String>,
) {
if let Some(tx) = &self.error_tx {
let report = InternalErrorReport::new(source, error, context);
if let Err(FibreTrySendError::Full(_report)) = tx.try_send(report) {
eprintln!("[fibre_logging:ERROR] Internal error channel full. Dropping error report.");
}
} else {
eprintln!("[fibre_logging:ERROR] {}: {}", source, error);
}
}
}