use camel_api::{CamelError, Exchange, IdentityProcessor, Value, ValueSource};
use std::future::Future;
use std::pin::Pin;
use std::task::{Context, Poll};
use tower::Service;
use tracing::{debug, info, trace, warn};
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum LogLevel {
Trace,
Debug,
Info,
Warn,
Error,
}
#[derive(Clone)]
pub struct LogProcessor {
inner: IdentityProcessor,
level: LogLevel,
message: String,
}
impl LogProcessor {
pub fn new(level: LogLevel, message: String) -> Self {
Self {
inner: IdentityProcessor,
level,
message,
}
}
}
impl Service<Exchange> for LogProcessor {
type Response = Exchange;
type Error = CamelError;
type Future = Pin<Box<dyn Future<Output = Result<Exchange, CamelError>> + Send>>;
fn poll_ready(&mut self, cx: &mut Context<'_>) -> Poll<Result<(), Self::Error>> {
self.inner.poll_ready(cx)
}
fn call(&mut self, exchange: Exchange) -> Self::Future {
let msg = self.message.clone();
let exchange_id = exchange.correlation_id.clone();
let body_preview = sanitize_preview(
&exchange
.input
.body
.as_text()
.unwrap_or("")
.chars()
.take(64)
.collect::<String>(),
);
debug!(exchange_id = %exchange_id, body_preview = %body_preview, "LogProcessor processing exchange");
match self.level {
LogLevel::Trace => trace!(exchange_id = %exchange_id, "{}", msg),
LogLevel::Debug => debug!(exchange_id = %exchange_id, "{}", msg),
LogLevel::Info => info!(exchange_id = %exchange_id, "{}", msg),
LogLevel::Warn => warn!(exchange_id = %exchange_id, "{}", msg),
LogLevel::Error => warn!(exchange_id = %exchange_id, "{}", msg),
}
self.inner.call(exchange)
}
}
#[derive(Clone)]
pub struct DynamicLog {
inner: IdentityProcessor,
level: LogLevel,
source: ValueSource,
}
impl DynamicLog {
pub fn new(level: LogLevel, source: impl Into<ValueSource>) -> Self {
Self {
inner: IdentityProcessor,
level,
source: source.into(),
}
}
}
fn log_value_to_string(value: Value) -> String {
match value {
Value::String(s) => s,
other => other.to_string(),
}
}
impl Service<Exchange> for DynamicLog {
type Response = Exchange;
type Error = CamelError;
type Future = Pin<Box<dyn Future<Output = Result<Exchange, CamelError>> + Send>>;
fn poll_ready(&mut self, cx: &mut Context<'_>) -> Poll<Result<(), Self::Error>> {
self.inner.poll_ready(cx)
}
fn call(&mut self, exchange: Exchange) -> Self::Future {
let source = self.source.clone();
let level = self.level;
Box::pin(async move {
let value = source.evaluate(&exchange).await?;
let exchange_id = exchange.correlation_id.clone();
let msg = sanitize_preview(&log_value_to_string(value));
match level {
LogLevel::Trace => trace!(exchange_id = %exchange_id, "{}", msg),
LogLevel::Debug => debug!(exchange_id = %exchange_id, "{}", msg),
LogLevel::Info => info!(exchange_id = %exchange_id, "{}", msg),
LogLevel::Warn => warn!(exchange_id = %exchange_id, "{}", msg),
LogLevel::Error => warn!(exchange_id = %exchange_id, "{}", msg),
}
Ok(exchange)
})
}
}
fn sanitize_preview(s: &str) -> String {
s.chars()
.map(|c| if c.is_control() { '\u{FFFD}' } else { c })
.collect()
}
#[cfg(test)]
mod tests {
use super::*;
use camel_api::body::Body;
use camel_api::{Message, Value};
use tower::ServiceExt;
#[tokio::test]
async fn test_log_processor_passes_exchange_through() {
let mut processor = LogProcessor::new(LogLevel::Info, "test message".into());
let exchange = Exchange::default();
let result = processor.call(exchange).await;
assert!(result.is_ok());
}
#[tokio::test]
async fn test_log_processor_preserves_exchange_body() {
let mut processor = LogProcessor::new(LogLevel::Debug, "debug message".into());
let mut exchange = Exchange::default();
exchange.input.body = Body::Text("test body".into());
let result = processor.call(exchange).await.unwrap();
assert_eq!(result.input.body.as_text(), Some("test body"));
}
fn sync_source<F>(f: F) -> camel_api::ValueSource
where
F: Fn(&Exchange) -> camel_api::Value + Send + Sync + 'static,
{
camel_api::ValueSource::Sync(std::sync::Arc::new(f))
}
#[tokio::test]
async fn test_dynamic_log_evaluates_body() {
let svc = DynamicLog::new(
LogLevel::Info,
sync_source(|ex: &Exchange| {
camel_api::Value::String(format!("body={}", ex.input.body.as_text().unwrap_or("")))
}),
);
let exchange = Exchange::new(Message::new("hello"));
let result = svc.oneshot(exchange).await.unwrap();
assert_eq!(result.input.body.as_text(), Some("hello"));
}
#[tokio::test]
async fn test_dynamic_log_evaluates_header() {
let svc = DynamicLog::new(
LogLevel::Info,
sync_source(|ex: &Exchange| {
let counter = ex
.input
.header("CamelTimerCounter")
.and_then(|v| v.as_i64())
.unwrap_or(0);
camel_api::Value::String(format!("{} World", counter))
}),
);
let mut msg = Message::new("");
msg.set_header("CamelTimerCounter", Value::Number(42.into()));
let exchange = Exchange::new(msg);
let result = svc.oneshot(exchange).await.unwrap();
assert_eq!(
result.input.header("CamelTimerCounter"),
Some(&Value::Number(42.into()))
);
}
#[tokio::test]
async fn dynamic_log_error_fails_step() {
use std::sync::Arc;
use camel_api::{BoxValueFuture, ValueSource};
let source = ValueSource::Async(Arc::new(|_: &Exchange| {
Box::pin(async { Err(CamelError::ProcessorError("log boom".into())) }) as BoxValueFuture
}));
let svc = DynamicLog::new(LogLevel::Info, source);
let result = svc.oneshot(Exchange::new(Message::new("x"))).await;
assert!(
result.is_err(),
"a failed log-message expression must fail the step, not log a null message"
);
}
#[test]
fn body_preview_strips_control_chars() {
let body_with_injection = "no\ttab\nnl\rcr\0null\u{1B}esc";
let preview: String = body_with_injection.chars().take(64).collect();
let sanitized = sanitize_preview(&preview);
assert!(!sanitized.contains('\n'), "newlines must be replaced");
assert!(
!sanitized.contains('\r'),
"carriage returns must be replaced"
);
assert!(!sanitized.contains('\t'), "tabs must be replaced");
assert!(!sanitized.contains('\0'), "null must be replaced");
assert!(
!sanitized.contains('\u{1B}'),
"ESC (ANSI injection) must be replaced"
);
assert!(
sanitized.contains('\u{FFFD}'),
"control chars must be replaced with U+FFFD, not deleted"
);
}
#[tokio::test]
async fn dynamic_log_sanitizes_control_chars() {
let svc = DynamicLog::new(
LogLevel::Info,
sync_source(|_ex: &Exchange| {
camel_api::Value::String("fake\nline\rinjection".to_string())
}),
);
let exchange = Exchange::default();
let result = svc.oneshot(exchange).await;
assert!(result.is_ok());
}
}