1use camel_api::{CamelError, Exchange, IdentityProcessor, Value, ValueSource};
2use std::future::Future;
3use std::pin::Pin;
4use std::task::{Context, Poll};
5use tower::Service;
6use tracing::{debug, info, trace, warn};
7
8#[derive(Debug, Clone, Copy, PartialEq, Eq)]
9pub enum LogLevel {
10 Trace,
11 Debug,
12 Info,
13 Warn,
14 Error,
15}
16
17#[derive(Clone)]
18pub struct LogProcessor {
19 inner: IdentityProcessor,
20 level: LogLevel,
21 message: String,
22}
23
24impl LogProcessor {
29 pub fn new(level: LogLevel, message: String) -> Self {
30 Self {
31 inner: IdentityProcessor,
32 level,
33 message,
34 }
35 }
36}
37
38impl Service<Exchange> for LogProcessor {
39 type Response = Exchange;
40 type Error = CamelError;
41 type Future = Pin<Box<dyn Future<Output = Result<Exchange, CamelError>> + Send>>;
42
43 fn poll_ready(&mut self, cx: &mut Context<'_>) -> Poll<Result<(), Self::Error>> {
44 self.inner.poll_ready(cx)
45 }
46
47 fn call(&mut self, exchange: Exchange) -> Self::Future {
48 let msg = self.message.clone();
49 let exchange_id = exchange.correlation_id.clone();
50 let body_preview = sanitize_preview(
51 &exchange
52 .input
53 .body
54 .as_text()
55 .unwrap_or("")
56 .chars()
57 .take(64)
58 .collect::<String>(),
59 );
60 debug!(exchange_id = %exchange_id, body_preview = %body_preview, "LogProcessor processing exchange");
61 match self.level {
62 LogLevel::Trace => trace!(exchange_id = %exchange_id, "{}", msg),
63 LogLevel::Debug => debug!(exchange_id = %exchange_id, "{}", msg),
64 LogLevel::Info => info!(exchange_id = %exchange_id, "{}", msg),
65 LogLevel::Warn => warn!(exchange_id = %exchange_id, "{}", msg),
66 LogLevel::Error => warn!(exchange_id = %exchange_id, "{}", msg),
68 }
69 self.inner.call(exchange)
70 }
71}
72
73#[derive(Clone)]
79pub struct DynamicLog {
80 inner: IdentityProcessor,
81 level: LogLevel,
82 source: ValueSource,
83}
84
85impl DynamicLog {
86 pub fn new(level: LogLevel, source: impl Into<ValueSource>) -> Self {
87 Self {
88 inner: IdentityProcessor,
89 level,
90 source: source.into(),
91 }
92 }
93}
94
95fn log_value_to_string(value: Value) -> String {
98 match value {
99 Value::String(s) => s,
100 other => other.to_string(),
101 }
102}
103
104impl Service<Exchange> for DynamicLog {
105 type Response = Exchange;
106 type Error = CamelError;
107 type Future = Pin<Box<dyn Future<Output = Result<Exchange, CamelError>> + Send>>;
108
109 fn poll_ready(&mut self, cx: &mut Context<'_>) -> Poll<Result<(), Self::Error>> {
110 self.inner.poll_ready(cx)
111 }
112
113 fn call(&mut self, exchange: Exchange) -> Self::Future {
114 let source = self.source.clone();
115 let level = self.level;
116 Box::pin(async move {
117 let value = source.evaluate(&exchange).await?;
118 let exchange_id = exchange.correlation_id.clone();
119 let msg = sanitize_preview(&log_value_to_string(value));
120 match level {
121 LogLevel::Trace => trace!(exchange_id = %exchange_id, "{}", msg),
122 LogLevel::Debug => debug!(exchange_id = %exchange_id, "{}", msg),
123 LogLevel::Info => info!(exchange_id = %exchange_id, "{}", msg),
124 LogLevel::Warn => warn!(exchange_id = %exchange_id, "{}", msg),
125 LogLevel::Error => warn!(exchange_id = %exchange_id, "{}", msg),
127 }
128 Ok(exchange)
129 })
130 }
131}
132
133fn sanitize_preview(s: &str) -> String {
137 s.chars()
138 .map(|c| if c.is_control() { '\u{FFFD}' } else { c })
139 .collect()
140}
141
142#[cfg(test)]
143mod tests {
144 use super::*;
145 use camel_api::body::Body;
146 use camel_api::{Message, Value};
147 use tower::ServiceExt;
148
149 #[tokio::test]
150 async fn test_log_processor_passes_exchange_through() {
151 let mut processor = LogProcessor::new(LogLevel::Info, "test message".into());
152 let exchange = Exchange::default();
153 let result = processor.call(exchange).await;
154 assert!(result.is_ok());
155 }
156
157 #[tokio::test]
158 async fn test_log_processor_preserves_exchange_body() {
159 let mut processor = LogProcessor::new(LogLevel::Debug, "debug message".into());
160 let mut exchange = Exchange::default();
161 exchange.input.body = Body::Text("test body".into());
162 let result = processor.call(exchange).await.unwrap();
163 assert_eq!(result.input.body.as_text(), Some("test body"));
164 }
165
166 fn sync_source<F>(f: F) -> camel_api::ValueSource
167 where
168 F: Fn(&Exchange) -> camel_api::Value + Send + Sync + 'static,
169 {
170 camel_api::ValueSource::Sync(std::sync::Arc::new(f))
171 }
172
173 #[tokio::test]
174 async fn test_dynamic_log_evaluates_body() {
175 let svc = DynamicLog::new(
176 LogLevel::Info,
177 sync_source(|ex: &Exchange| {
178 camel_api::Value::String(format!("body={}", ex.input.body.as_text().unwrap_or("")))
179 }),
180 );
181 let exchange = Exchange::new(Message::new("hello"));
182 let result = svc.oneshot(exchange).await.unwrap();
183 assert_eq!(result.input.body.as_text(), Some("hello"));
185 }
186
187 #[tokio::test]
188 async fn test_dynamic_log_evaluates_header() {
189 let svc = DynamicLog::new(
190 LogLevel::Info,
191 sync_source(|ex: &Exchange| {
192 let counter = ex
193 .input
194 .header("CamelTimerCounter")
195 .and_then(|v| v.as_i64())
196 .unwrap_or(0);
197 camel_api::Value::String(format!("{} World", counter))
198 }),
199 );
200 let mut msg = Message::new("");
201 msg.set_header("CamelTimerCounter", Value::Number(42.into()));
202 let exchange = Exchange::new(msg);
203 let result = svc.oneshot(exchange).await.unwrap();
204 assert_eq!(
206 result.input.header("CamelTimerCounter"),
207 Some(&Value::Number(42.into()))
208 );
209 }
210
211 #[tokio::test]
212 async fn dynamic_log_error_fails_step() {
213 use std::sync::Arc;
214
215 use camel_api::{BoxValueFuture, ValueSource};
216
217 let source = ValueSource::Async(Arc::new(|_: &Exchange| {
218 Box::pin(async { Err(CamelError::ProcessorError("log boom".into())) }) as BoxValueFuture
219 }));
220 let svc = DynamicLog::new(LogLevel::Info, source);
221
222 let result = svc.oneshot(Exchange::new(Message::new("x"))).await;
223 assert!(
224 result.is_err(),
225 "a failed log-message expression must fail the step, not log a null message"
226 );
227 }
228
229 #[test]
230 fn body_preview_strips_control_chars() {
231 let body_with_injection = "no\ttab\nnl\rcr\0null\u{1B}esc";
232 let preview: String = body_with_injection.chars().take(64).collect();
233 let sanitized = sanitize_preview(&preview);
234 assert!(!sanitized.contains('\n'), "newlines must be replaced");
235 assert!(
236 !sanitized.contains('\r'),
237 "carriage returns must be replaced"
238 );
239 assert!(!sanitized.contains('\t'), "tabs must be replaced");
240 assert!(!sanitized.contains('\0'), "null must be replaced");
241 assert!(
242 !sanitized.contains('\u{1B}'),
243 "ESC (ANSI injection) must be replaced"
244 );
245 assert!(
247 sanitized.contains('\u{FFFD}'),
248 "control chars must be replaced with U+FFFD, not deleted"
249 );
250 }
251
252 #[tokio::test]
253 async fn dynamic_log_sanitizes_control_chars() {
254 let svc = DynamicLog::new(
255 LogLevel::Info,
256 sync_source(|_ex: &Exchange| {
257 camel_api::Value::String("fake\nline\rinjection".to_string())
258 }),
259 );
260 let exchange = Exchange::default();
264 let result = svc.oneshot(exchange).await;
265 assert!(result.is_ok());
266 }
267}