Skip to main content

camel_processor/
log.rs

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
24// TODO(PROC-004): Add metrics instrumentation — processed count, error count, latency histograms
25// are not yet instrumented on processors. Consider wiring MetricsCollector into LogProcessor and
26// incrementing a counter on each call, recording elapsed time, and tracking errors.
27
28impl 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            // log-policy: handler-owned
67            LogLevel::Error => warn!(exchange_id = %exchange_id, "{}", msg),
68        }
69        self.inner.call(exchange)
70    }
71}
72
73/// A log processor that evaluates a message expression against the Exchange at call-time.
74/// Analogous to [`DynamicSetHeader`](crate::dynamic_set_header::DynamicSetHeader).
75///
76/// A failed evaluation fails the step BEFORE logging — no log record is
77/// emitted with a null message.
78#[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
95/// Render an evaluated log-message value: strings pass through unquoted,
96/// everything else uses its JSON representation.
97fn 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                // log-policy: handler-owned
126                LogLevel::Error => warn!(exchange_id = %exchange_id, "{}", msg),
127            }
128            Ok(exchange)
129        })
130    }
131}
132
133/// Strip control characters that could enable log injection.
134/// Replaces control chars with U+FFFD (replacement char) to preserve
135/// alignment while making injection visible.
136fn 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        // Exchange passes through unchanged
184        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        // Exchange passes through unchanged
205        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        // U+FFFD replacement char must appear (replacement, not deletion)
246        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        // The exchange passes through; sanitization happens inside call().
261        // We verify the exchange is unchanged and the call succeeds — the
262        // sanitization itself is validated by body_preview_strips_control_chars.
263        let exchange = Exchange::default();
264        let result = svc.oneshot(exchange).await;
265        assert!(result.is_ok());
266    }
267}