Skip to main content

camel_component_log/
lib.rs

1//! Log component for rust-camel — routes exchange body and headers to the
2//! `tracing` subsystem at a configurable log level, primarily for debugging.
3//!
4//! Main types: `LogComponent`, `LogEndpoint`, `LogProducer`, `LogLevel`.
5
6use std::future::Future;
7use std::pin::Pin;
8use std::str::FromStr;
9use std::sync::Arc;
10use std::sync::atomic::{AtomicUsize, Ordering};
11use std::task::{Context, Poll};
12
13use tower::Service;
14use tracing::{debug, error, info, trace, warn};
15
16use camel_component_api::UriConfig;
17use camel_component_api::parse_uri;
18use camel_component_api::{BoxProcessor, CamelError, Exchange};
19use camel_component_api::{
20    Component, ComponentMetadata, Consumer, Endpoint, ProducerContext, RuntimeObservability,
21};
22
23// ---------------------------------------------------------------------------
24// LogLevel
25// ---------------------------------------------------------------------------
26
27/// Log level for the log component.
28#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
29pub enum LogLevel {
30    Trace,
31    Debug,
32    #[default]
33    Info,
34    Warn,
35    Error,
36}
37
38impl FromStr for LogLevel {
39    type Err = String;
40
41    fn from_str(s: &str) -> Result<Self, Self::Err> {
42        parse_log_level(s).map_err(|e| e.to_string())
43    }
44}
45
46fn parse_log_level(s: &str) -> Result<LogLevel, CamelError> {
47    match s.to_uppercase().as_str() {
48        "TRACE" => Ok(LogLevel::Trace),
49        "DEBUG" => Ok(LogLevel::Debug),
50        "INFO" => Ok(LogLevel::Info),
51        "WARN" | "WARNING" => Ok(LogLevel::Warn),
52        "ERROR" => Ok(LogLevel::Error),
53        _ => Err(CamelError::Config(format!(
54            "unknown log level: '{}'. Valid: TRACE, DEBUG, INFO, WARN, ERROR",
55            s
56        ))),
57    }
58}
59
60// ---------------------------------------------------------------------------
61// LogConfig
62// ---------------------------------------------------------------------------
63
64/// Configuration parsed from a log URI.
65///
66/// Format: `log:category?level=info&showHeaders=true&showBody=true&logMask=false&showStreamInfo=false&groupSize=10`
67#[derive(Debug, Clone, UriConfig)]
68#[uri_scheme = "log"]
69#[uri_config(
70    skip_impl,
71    metadata(
72        scheme = "log",
73        description = "Logs exchanges at a configurable level",
74        producer
75    ),
76    crate = "camel_component_api"
77)]
78pub struct LogConfig {
79    /// Log category (the path portion of the URI).
80    pub category: String,
81    /// Log level. Default: Info.
82    #[uri_param(
83        kind = "enum:TRACE,DEBUG,INFO,WARN,WARNING,ERROR",
84        default = "INFO",
85        desc = "Log level"
86    )]
87    pub level: LogLevel,
88    /// Whether to include headers in the log output.
89    #[uri_param(
90        name = "showHeaders",
91        default = "false",
92        desc = "Include headers in log output"
93    )]
94    pub show_headers: bool,
95    /// Whether to include the body in the log output.
96    #[uri_param(
97        name = "showBody",
98        default = "true",
99        desc = "Include body in log output"
100    )]
101    pub show_body: bool,
102    /// Maximum number of characters for the body in log output.
103    /// Bodies longer than this are truncated. `None` means no limit.
104    #[uri_param(
105        name = "maxChars",
106        desc = "Truncate body to N characters. Omit for no limit"
107    )]
108    pub max_chars: Option<usize>,
109    /// When true, redact sensitive headers and body in log output.
110    /// Headers matching `/(?i)(password|secret|token|key|auth|credential)/` → `[REDACTED]`.
111    /// Body → `[Body redacted by logMask]`.
112    #[uri_param(
113        name = "logMask",
114        default = "false",
115        desc = "Redact sensitive headers and body"
116    )]
117    pub log_mask: bool,
118    /// When true, show stream origin metadata. When false (default), show `[Stream]`.
119    #[uri_param(
120        name = "showStreamInfo",
121        default = "false",
122        desc = "Show stream origin info"
123    )]
124    pub show_stream_info: bool,
125    /// When set, only emit a log every `n` exchanges (group logging).
126    /// The log message includes the exchange count.
127    #[uri_param(
128        name = "groupSize",
129        desc = "Emit log every N exchanges for group logging"
130    )]
131    pub group_size: Option<usize>,
132}
133
134impl UriConfig for LogConfig {
135    fn scheme() -> &'static str {
136        "log"
137    }
138
139    fn from_uri(uri: &str) -> Result<Self, CamelError> {
140        let parts = parse_uri(uri)?;
141        Self::from_components(parts)
142    }
143
144    fn from_components(parts: camel_component_api::UriComponents) -> Result<Self, CamelError> {
145        if parts.scheme != Self::scheme() {
146            return Err(CamelError::InvalidUri(format!(
147                "expected scheme '{}' but got '{}'",
148                Self::scheme(),
149                parts.scheme
150            )));
151        }
152
153        let level = match parts.params.get("level") {
154            Some(raw) => parse_log_level(raw)?,
155            None => LogLevel::Info,
156        };
157
158        let show_headers = match parts.params.get("showHeaders") {
159            Some(raw) => raw.parse::<bool>().map_err(|_| {
160                CamelError::InvalidUri(format!("invalid boolean value for showHeaders: {raw}"))
161            })?,
162            None => false,
163        };
164
165        let show_body = match parts.params.get("showBody") {
166            Some(raw) => raw.parse::<bool>().map_err(|_| {
167                CamelError::InvalidUri(format!("invalid boolean value for showBody: {raw}"))
168            })?,
169            None => true,
170        };
171
172        let max_chars = match parts.params.get("maxChars") {
173            Some(raw) => Some(raw.parse::<usize>().map_err(|_| {
174                CamelError::InvalidUri(format!("invalid integer value for maxChars: {raw}"))
175            })?),
176            None => None,
177        };
178
179        let log_mask = match parts.params.get("logMask") {
180            Some(raw) => raw.parse::<bool>().map_err(|_| {
181                CamelError::InvalidUri(format!("invalid boolean value for logMask: {raw}"))
182            })?,
183            None => false,
184        };
185
186        let show_stream_info = match parts.params.get("showStreamInfo") {
187            Some(raw) => raw.parse::<bool>().map_err(|_| {
188                CamelError::InvalidUri(format!("invalid boolean value for showStreamInfo: {raw}"))
189            })?,
190            None => false,
191        };
192
193        let group_size = match parts.params.get("groupSize") {
194            Some(raw) => Some(raw.parse::<usize>().map_err(|_| {
195                CamelError::InvalidUri(format!("invalid integer value for groupSize: {raw}"))
196            })?),
197            None => None,
198        };
199
200        Ok(Self {
201            category: parts.path,
202            level,
203            show_headers,
204            show_body,
205            max_chars,
206            log_mask,
207            show_stream_info,
208            group_size,
209        })
210    }
211}
212
213// ---------------------------------------------------------------------------
214// LogComponent
215// ---------------------------------------------------------------------------
216
217/// The Log component logs exchange information using `tracing`.
218pub struct LogComponent;
219
220impl LogComponent {
221    pub fn new() -> Self {
222        Self
223    }
224}
225
226impl Default for LogComponent {
227    fn default() -> Self {
228        Self::new()
229    }
230}
231
232impl Component for LogComponent {
233    fn scheme(&self) -> &str {
234        "log"
235    }
236
237    fn metadata(&self) -> ComponentMetadata {
238        LogConfig::metadata()
239    }
240
241    fn create_endpoint(
242        &self,
243        uri: &str,
244        _ctx: &dyn camel_component_api::ComponentContext,
245    ) -> Result<Box<dyn Endpoint>, CamelError> {
246        let config = LogConfig::from_uri(uri)?;
247        Ok(Box::new(LogEndpoint {
248            uri: uri.to_string(),
249            config,
250        }))
251    }
252}
253
254// ---------------------------------------------------------------------------
255// LogEndpoint
256// ---------------------------------------------------------------------------
257
258struct LogEndpoint {
259    uri: String,
260    config: LogConfig,
261}
262
263impl Endpoint for LogEndpoint {
264    fn uri(&self) -> &str {
265        &self.uri
266    }
267
268    fn create_consumer(
269        &self,
270        _rt: Arc<dyn RuntimeObservability>,
271    ) -> Result<Box<dyn Consumer>, CamelError> {
272        Err(CamelError::EndpointCreationFailed(
273            "log endpoint does not support consumers".to_string(),
274        ))
275    }
276
277    fn create_producer(
278        &self,
279        _rt: Arc<dyn RuntimeObservability>,
280        _ctx: &ProducerContext,
281    ) -> Result<BoxProcessor, CamelError> {
282        Ok(BoxProcessor::new(LogProducer::new(self.config.clone())))
283    }
284}
285
286// ---------------------------------------------------------------------------
287// LogProducer
288// ---------------------------------------------------------------------------
289
290#[derive(Clone)]
291struct LogProducer {
292    config: LogConfig,
293    exchange_count: Arc<AtomicUsize>,
294}
295
296impl LogProducer {
297    fn new(config: LogConfig) -> Self {
298        Self {
299            config,
300            exchange_count: Arc::new(AtomicUsize::new(0)),
301        }
302    }
303
304    /// Returns true if the header key matches sensitive patterns.
305    fn is_sensitive_header(key: &str) -> bool {
306        let lower = key.to_lowercase();
307        let sensitive_keywords = [
308            "password",
309            "passwd",
310            "secret",
311            "token",
312            "apikey",
313            "api-key",
314            "api_key",
315            "authorization",
316            "auth",
317            "credential",
318            "private",
319            "signature",
320        ];
321        sensitive_keywords.iter().any(|kw| {
322            lower == *kw
323                || lower.ends_with(&format!("-{kw}"))
324                || lower.ends_with(&format!("_{kw}"))
325                || lower.starts_with(&format!("{kw}-"))
326                || lower.starts_with(&format!("{kw}_"))
327        })
328    }
329
330    fn format_exchange(&self, exchange: &Exchange, count: usize) -> String {
331        let mut parts = Vec::new();
332
333        if self.config.show_body {
334            let body_str = if self.config.log_mask {
335                "[Body redacted by logMask]".to_string()
336            } else {
337                match &exchange.input.body {
338                    camel_component_api::Body::Empty => "[empty]".to_string(),
339                    camel_component_api::Body::Text(s) => s.clone(),
340                    camel_component_api::Body::Json(v) => v.to_string(),
341                    camel_component_api::Body::Xml(s) => s.clone(),
342                    camel_component_api::Body::Bytes(b) => format!("[{} bytes]", b.len()),
343                    camel_component_api::Body::Stream(s) => {
344                        if self.config.show_stream_info {
345                            format!("[Stream: origin={:?}]", s.metadata.origin)
346                        } else {
347                            "[Stream]".to_string()
348                        }
349                    }
350                    // Future body variants render as an opaque placeholder.
351                    _ => "[unsupported body type]".to_string(),
352                }
353            };
354
355            let body_str = if let Some(limit) = self.config.max_chars {
356                body_str.chars().take(limit).collect()
357            } else {
358                body_str
359            };
360
361            parts.push(format!("Body: {body_str}"));
362        }
363
364        if self.config.show_headers && !exchange.input.headers.is_empty() {
365            let headers: Vec<String> = exchange
366                .input
367                .headers
368                .iter()
369                .map(|(k, v)| {
370                    if self.config.log_mask && Self::is_sensitive_header(k) {
371                        format!("{k}=[REDACTED]")
372                    } else {
373                        format!("{k}={v}")
374                    }
375                })
376                .collect();
377            parts.push(format!("Headers: {{{}}}", headers.join(", ")));
378        }
379
380        if parts.is_empty() {
381            format!("[{}] Exchange received", self.config.category)
382        } else if self.config.group_size.is_some() {
383            format!(
384                "[{}] Group of {count}: {}",
385                self.config.category,
386                parts.join(" | ")
387            )
388        } else {
389            format!("[{}] {}", self.config.category, parts.join(" | "))
390        }
391    }
392}
393
394impl Service<Exchange> for LogProducer {
395    type Response = Exchange;
396    type Error = CamelError;
397    type Future = Pin<Box<dyn Future<Output = Result<Exchange, CamelError>> + Send>>;
398
399    fn poll_ready(&mut self, _cx: &mut Context<'_>) -> Poll<Result<(), Self::Error>> {
400        Poll::Ready(Ok(()))
401    }
402
403    fn call(&mut self, exchange: Exchange) -> Self::Future {
404        let count = self.exchange_count.fetch_add(1, Ordering::Relaxed) + 1;
405
406        // Group logging: only emit every group_size exchanges
407        if !Self::should_emit(count, self.config.group_size) {
408            return Box::pin(async move { Ok(exchange) });
409        }
410
411        let msg = self.format_exchange(&exchange, count);
412        let level = self.config.level;
413
414        Box::pin(async move {
415            match level {
416                LogLevel::Trace => trace!("{msg}"),
417                LogLevel::Debug => debug!("{msg}"),
418                LogLevel::Info => info!("{msg}"),
419                LogLevel::Warn => warn!("{msg}"),
420                // Deliberately unannotated: LogProducer is a user-output
421                // mechanism — the route selects this level — so the call sits
422                // outside ADR-0012's operational convention (see this crate's
423                // CONTEXT.md). The lint-log-levels exclusion is symbol-bound
424                // (structural, in scripts/xtask), and
425                // tests/logproducer_exclusion_regression.rs guards that this
426                // crate carries zero level-policy annotations.
427                LogLevel::Error => error!("{msg}"),
428            }
429
430            Ok(exchange)
431        })
432    }
433}
434
435impl LogProducer {
436    /// Decide whether the `count`-th exchange should emit a log line.
437    /// Without grouping every exchange emits; with grouping only multiples of
438    /// `group_size` emit (counts in between pass through silently). Extracted
439    /// from `call()` so the gating is unit-testable (rc-h0bv).
440    fn should_emit(count: usize, group_size: Option<usize>) -> bool {
441        match group_size {
442            None => true,
443            Some(gs) => count.is_multiple_of(gs),
444        }
445    }
446}
447
448// ---------------------------------------------------------------------------
449// Tests
450// ---------------------------------------------------------------------------
451
452#[cfg(test)]
453mod tests {
454    use camel_component_api::test_support::PanicRuntimeObservability;
455    fn rt() -> std::sync::Arc<dyn camel_component_api::RuntimeObservability> {
456        std::sync::Arc::new(PanicRuntimeObservability)
457    }
458    use super::*;
459    use camel_component_api::Body;
460    use camel_component_api::Message;
461    use camel_component_api::NoOpComponentContext;
462    use tower::ServiceExt;
463
464    fn test_producer_ctx() -> ProducerContext {
465        ProducerContext::new()
466    }
467
468    #[test]
469    fn test_log_config_defaults() {
470        let config = LogConfig::from_uri("log:myCategory").unwrap();
471        assert_eq!(config.category, "myCategory");
472        assert_eq!(config.level, LogLevel::Info);
473        assert!(!config.show_headers);
474        assert!(config.show_body);
475    }
476
477    #[test]
478    fn test_log_config_with_params() {
479        let config =
480            LogConfig::from_uri("log:app?level=debug&showHeaders=true&showBody=false").unwrap();
481        assert_eq!(config.category, "app");
482        assert_eq!(config.level, LogLevel::Debug);
483        assert!(config.show_headers);
484        assert!(!config.show_body);
485    }
486
487    #[test]
488    fn test_log_config_wrong_scheme() {
489        let result = LogConfig::from_uri("timer:tick");
490        assert!(result.is_err());
491    }
492
493    #[test]
494    fn test_log_component_scheme() {
495        let component = LogComponent::new();
496        assert_eq!(component.scheme(), "log");
497    }
498
499    #[test]
500    fn test_log_component_default() {
501        let component = LogComponent;
502        assert_eq!(component.scheme(), "log");
503    }
504
505    #[test]
506    fn test_log_level_from_str_variants() {
507        assert_eq!("trace".parse::<LogLevel>().unwrap(), LogLevel::Trace);
508        assert_eq!("DEBUG".parse::<LogLevel>().unwrap(), LogLevel::Debug);
509        assert_eq!("Info".parse::<LogLevel>().unwrap(), LogLevel::Info);
510        assert_eq!("warning".parse::<LogLevel>().unwrap(), LogLevel::Warn);
511        assert_eq!("error".parse::<LogLevel>().unwrap(), LogLevel::Error);
512    }
513
514    #[test]
515    fn test_log_level_from_str_invalid() {
516        let err = "nope".parse::<LogLevel>().unwrap_err();
517        assert_eq!(
518            err,
519            "Configuration error: unknown log level: 'nope'. Valid: TRACE, DEBUG, INFO, WARN, ERROR"
520        );
521    }
522
523    #[test]
524    fn test_log_config_invalid_level_rejected() {
525        let err = LogConfig::from_uri("log:test?level=invalid").unwrap_err();
526        assert!(
527            err.to_string()
528                .contains("unknown log level: 'invalid'. Valid: TRACE, DEBUG, INFO, WARN, ERROR")
529        );
530    }
531
532    #[test]
533    fn test_valid_log_levels_accepted() {
534        assert!(parse_log_level("DEBUG").is_ok());
535        assert!(parse_log_level("info").is_ok());
536        assert!(parse_log_level("WARN").is_ok());
537        assert!(parse_log_level("WARNING").is_ok());
538    }
539
540    #[test]
541    fn test_invalid_log_level_rejected() {
542        assert!(parse_log_level("VERBOSE").is_err());
543        assert!(parse_log_level("").is_err());
544        assert!(parse_log_level("log").is_err());
545    }
546
547    #[test]
548    fn test_metadata_level_enum_covers_runtime_vocab() {
549        // The lint's R-URI-known kind check validates `level` values against
550        // the metadata Enum variants; the runtime accepts the alias WARNING
551        // (any case), so the variant list must cover it or camel lint false-
552        // positives on runtime-legal routes (rc-68q6 review finding).
553        use camel_component_api::{ComponentMetadata, OptionKind};
554        let meta: ComponentMetadata = LogConfig::metadata();
555        let level = meta
556            .uri_options
557            .iter()
558            .find(|o| o.name == "level")
559            .expect("level option must exist");
560        let OptionKind::Enum(variants) = &level.kind else {
561            panic!("level option kind must be Enum; got {:?}", level.kind);
562        };
563        let accepted = ["TRACE", "DEBUG", "INFO", "WARN", "WARNING", "ERROR"];
564        for value in accepted {
565            assert!(
566                variants.iter().any(|v| v.eq_ignore_ascii_case(value)),
567                "metadata level variants must cover runtime-accepted `{value}`; got {variants:?}"
568            );
569        }
570    }
571
572    #[test]
573    fn test_log_endpoint_uri() {
574        let component = LogComponent::new();
575        let endpoint = component
576            .create_endpoint("log:uri-check", &NoOpComponentContext)
577            .unwrap();
578        assert_eq!(endpoint.uri(), "log:uri-check");
579    }
580
581    #[test]
582    fn test_log_endpoint_no_consumer() {
583        let component = LogComponent::new();
584        let endpoint = component
585            .create_endpoint("log:info", &NoOpComponentContext)
586            .unwrap();
587        assert!(endpoint.create_consumer(rt()).is_err());
588    }
589
590    #[test]
591    fn test_log_endpoint_creates_producer() {
592        let ctx = test_producer_ctx();
593        let component = LogComponent::new();
594        let endpoint = component
595            .create_endpoint("log:info", &NoOpComponentContext)
596            .unwrap();
597        assert!(endpoint.create_producer(rt(), &ctx).is_ok());
598    }
599
600    #[tokio::test]
601    async fn test_log_producer_processes_exchange() {
602        let ctx = test_producer_ctx();
603        let component = LogComponent::new();
604        let endpoint = component
605            .create_endpoint("log:test?showHeaders=true", &NoOpComponentContext)
606            .unwrap();
607        let producer = endpoint.create_producer(rt(), &ctx).unwrap();
608
609        let mut exchange = Exchange::new(Message::new("hello world"));
610        exchange
611            .input
612            .set_header("source", serde_json::Value::String("test".into()));
613
614        let result = producer.oneshot(exchange).await.unwrap();
615        // Log producer passes exchange through unchanged
616        assert_eq!(result.input.body.as_text(), Some("hello world"));
617    }
618
619    #[test]
620    fn test_format_exchange_without_body_or_headers() {
621        let producer = LogProducer::new(LogConfig {
622            category: "cat".to_string(),
623            level: LogLevel::Info,
624            show_headers: false,
625            show_body: false,
626            max_chars: None,
627            log_mask: false,
628            show_stream_info: false,
629            group_size: None,
630        });
631        let exchange = Exchange::new(Message::new("ignored"));
632        let formatted = producer.format_exchange(&exchange, 1);
633        assert_eq!(formatted, "[cat] Exchange received");
634    }
635
636    #[test]
637    fn test_format_exchange_body_variants() {
638        let base = LogProducer::new(LogConfig {
639            category: "cat".to_string(),
640            level: LogLevel::Info,
641            show_headers: false,
642            show_body: true,
643            max_chars: None,
644            log_mask: false,
645            show_stream_info: true,
646            group_size: None,
647        });
648
649        let empty = Exchange::new(Message::default());
650        assert!(base.format_exchange(&empty, 1).contains("Body: [empty]"));
651
652        let mut json_msg = Message::new("");
653        json_msg.body = Body::Json(serde_json::json!({"k":"v"}));
654        let json_ex = Exchange::new(json_msg);
655        assert!(
656            base.format_exchange(&json_ex, 2)
657                .contains("Body: {\"k\":\"v\"}")
658        );
659
660        let mut xml_msg = Message::new("");
661        xml_msg.body = Body::Xml("<a/>".to_string());
662        let xml_ex = Exchange::new(xml_msg);
663        assert!(base.format_exchange(&xml_ex, 3).contains("Body: <a/>"));
664
665        let mut bytes_msg = Message::new("");
666        bytes_msg.body = Body::Bytes(b"abc".to_vec().into());
667        let bytes_ex = Exchange::new(bytes_msg);
668        assert!(
669            base.format_exchange(&bytes_ex, 4)
670                .contains("Body: [3 bytes]")
671        );
672    }
673
674    #[test]
675    fn test_log_truncates_large_body() {
676        let producer = LogProducer::new(LogConfig {
677            category: "trunc".to_string(),
678            level: LogLevel::Info,
679            show_headers: false,
680            show_body: true,
681            max_chars: Some(10),
682            log_mask: false,
683            show_stream_info: false,
684            group_size: None,
685        });
686
687        let long_body = "a".repeat(100);
688        let exchange = Exchange::new(Message::new(long_body));
689        let formatted = producer.format_exchange(&exchange, 1);
690
691        // Extract the body part from "[trunc] Body: ..."
692        let body_part = formatted.split_once("Body: ").unwrap().1;
693        assert!(
694            body_part.chars().count() <= 10,
695            "expected body <= 10 chars, got {} chars: {body_part:?}",
696            body_part.chars().count()
697        );
698    }
699
700    #[test]
701    fn test_log_truncates_multibyte_body() {
702        let producer = LogProducer::new(LogConfig {
703            category: "mb".to_string(),
704            level: LogLevel::Info,
705            show_headers: false,
706            show_body: true,
707            max_chars: Some(3),
708            log_mask: false,
709            show_stream_info: false,
710            group_size: None,
711        });
712
713        let exchange = Exchange::new(Message::new("日本語测试"));
714        let formatted = producer.format_exchange(&exchange, 1);
715
716        let body_part = formatted.split_once("Body: ").unwrap().1;
717        assert_eq!(
718            body_part.chars().count(),
719            3,
720            "expected 3 chars, got {}: {body_part:?}",
721            body_part.chars().count()
722        );
723        // body_part is &str (from split_once), so UTF-8 validity is guaranteed
724        // by the type — no runtime check needed.
725    }
726
727    #[test]
728    fn test_log_truncates_multibyte_no_panic() {
729        let producer = LogProducer::new(LogConfig {
730            category: "cafe".to_string(),
731            level: LogLevel::Info,
732            show_headers: false,
733            show_body: true,
734            max_chars: Some(4),
735            log_mask: false,
736            show_stream_info: false,
737            group_size: None,
738        });
739
740        let exchange = Exchange::new(Message::new("café"));
741        let formatted = producer.format_exchange(&exchange, 1);
742
743        let body_part = formatted.split_once("Body: ").unwrap().1;
744        assert_eq!(
745            body_part.chars().count(),
746            4,
747            "expected 4 chars, got {}: {body_part:?}",
748            body_part.chars().count()
749        );
750    }
751
752    #[test]
753    fn test_log_no_truncation_when_max_chars_unset() {
754        let producer = LogProducer::new(LogConfig {
755            category: "notrunc".to_string(),
756            level: LogLevel::Info,
757            show_headers: false,
758            show_body: true,
759            max_chars: None,
760            log_mask: false,
761            show_stream_info: false,
762            group_size: None,
763        });
764
765        let long_body = "b".repeat(200);
766        let exchange = Exchange::new(Message::new(long_body));
767        let formatted = producer.format_exchange(&exchange, 1);
768
769        let body_part = formatted.split_once("Body: ").unwrap().1;
770        assert_eq!(body_part.len(), 200);
771    }
772
773    #[test]
774    fn test_log_config_max_chars_param() {
775        let config = LogConfig::from_uri("log:test?maxChars=50").unwrap();
776        assert_eq!(config.max_chars, Some(50));
777    }
778
779    #[test]
780    fn test_log_config_max_chars_default_unset() {
781        let config = LogConfig::from_uri("log:test").unwrap();
782        assert_eq!(config.max_chars, None);
783    }
784
785    // --- LOG-007: logMask tests ---
786
787    #[test]
788    fn test_log_config_log_mask_param() {
789        let config = LogConfig::from_uri("log:test?logMask=true").unwrap();
790        assert!(config.log_mask);
791    }
792
793    #[test]
794    fn test_log_config_log_mask_default_false() {
795        let config = LogConfig::from_uri("log:test").unwrap();
796        assert!(!config.log_mask);
797    }
798
799    #[test]
800    fn test_log_mask_redacts_sensitive_headers() {
801        let producer = LogProducer::new(LogConfig {
802            category: "cat".to_string(),
803            level: LogLevel::Info,
804            show_headers: true,
805            show_body: false,
806            max_chars: None,
807            log_mask: true,
808            show_stream_info: false,
809            group_size: None,
810        });
811
812        let mut exchange = Exchange::new(Message::new("body"));
813        exchange.input.set_header(
814            "X-Auth-Token",
815            serde_json::Value::String("secret123".into()),
816        );
817        exchange
818            .input
819            .set_header("password", serde_json::Value::String("hunter2".into()));
820        exchange
821            .input
822            .set_header("ApiKey", serde_json::Value::String("abc".into()));
823        exchange
824            .input
825            .set_header("normal-header", serde_json::Value::String("visible".into()));
826        exchange.input.set_header(
827            "user-credential",
828            serde_json::Value::String("sensitive".into()),
829        );
830        exchange
831            .input
832            .set_header("secret-value", serde_json::Value::String("hidden".into()));
833
834        let formatted = producer.format_exchange(&exchange, 1);
835        assert!(
836            formatted.contains("X-Auth-Token=[REDACTED]"),
837            "auth header must be redacted: {formatted}"
838        );
839        assert!(
840            formatted.contains("password=[REDACTED]"),
841            "password header must be redacted: {formatted}"
842        );
843        assert!(
844            formatted.contains("ApiKey=[REDACTED]"),
845            "key header must be redacted: {formatted}"
846        );
847        assert!(
848            formatted.contains("normal-header=\"visible\""),
849            "normal header must be visible: {formatted}"
850        );
851        assert!(
852            formatted.contains("user-credential=[REDACTED]"),
853            "credential header must be redacted: {formatted}"
854        );
855        assert!(
856            formatted.contains("secret-value=[REDACTED]"),
857            "secret header must be redacted: {formatted}"
858        );
859    }
860
861    #[test]
862    fn test_log_mask_redacts_body() {
863        let producer = LogProducer::new(LogConfig {
864            category: "cat".to_string(),
865            level: LogLevel::Info,
866            show_headers: false,
867            show_body: true,
868            max_chars: None,
869            log_mask: true,
870            show_stream_info: false,
871            group_size: None,
872        });
873
874        let exchange = Exchange::new(Message::new("sensitive body content"));
875        let formatted = producer.format_exchange(&exchange, 1);
876        assert!(
877            formatted.contains("[Body redacted by logMask]"),
878            "body must be redacted: {formatted}"
879        );
880        assert!(
881            !formatted.contains("sensitive body content"),
882            "body content must not appear: {formatted}"
883        );
884    }
885
886    #[test]
887    fn test_log_mask_off_shows_data() {
888        let producer = LogProducer::new(LogConfig {
889            category: "cat".to_string(),
890            level: LogLevel::Info,
891            show_headers: true,
892            show_body: true,
893            max_chars: None,
894            log_mask: false,
895            show_stream_info: false,
896            group_size: None,
897        });
898
899        let mut exchange = Exchange::new(Message::new("visible body"));
900        exchange
901            .input
902            .set_header("password", serde_json::Value::String("hunter2".into()));
903
904        let formatted = producer.format_exchange(&exchange, 1);
905        assert!(
906            formatted.contains("visible body"),
907            "body must be visible when mask off: {formatted}"
908        );
909        assert!(
910            formatted.contains("hunter2"),
911            "header value must be visible when mask off: {formatted}"
912        );
913    }
914
915    // --- LOG-002: showStreamInfo tests ---
916
917    #[test]
918    fn test_log_stream_show_info() {
919        // show_stream_info = false → just [Stream]
920        let producer_no_info = LogProducer::new(LogConfig {
921            category: "cat".to_string(),
922            level: LogLevel::Info,
923            show_headers: false,
924            show_body: true,
925            max_chars: None,
926            log_mask: false,
927            show_stream_info: false,
928            group_size: None,
929        });
930
931        let mut msg = Message::new("");
932        msg.body = Body::Stream(camel_component_api::StreamBody {
933            stream: std::sync::Arc::new(tokio::sync::Mutex::new(None)),
934            metadata: camel_component_api::StreamMetadata {
935                origin: Some("file:///data/test.txt".to_string()),
936                ..Default::default()
937            },
938        });
939        let exchange = Exchange::new(msg);
940        let formatted = producer_no_info.format_exchange(&exchange, 1);
941        assert!(
942            formatted.contains("Body: [Stream]"),
943            "must show [Stream] when show_stream_info=false: {formatted}"
944        );
945
946        // show_stream_info = true → [Stream: origin=...]
947        let producer_with_info = LogProducer::new(LogConfig {
948            category: "cat".to_string(),
949            level: LogLevel::Info,
950            show_headers: false,
951            show_body: true,
952            max_chars: None,
953            log_mask: false,
954            show_stream_info: true,
955            group_size: None,
956        });
957
958        let formatted = producer_with_info.format_exchange(&exchange, 1);
959        assert!(
960            formatted.contains("Body: [Stream: origin=Some(\"file:///data/test.txt\")]"),
961            "must show origin when show_stream_info=true: {formatted}"
962        );
963    }
964
965    // --- LOG-008: groupSize tests ---
966
967    #[test]
968    fn test_log_group_size() {
969        let producer = LogProducer::new(LogConfig {
970            category: "cat".to_string(),
971            level: LogLevel::Info,
972            show_headers: false,
973            show_body: true,
974            max_chars: None,
975            log_mask: false,
976            show_stream_info: false,
977            group_size: Some(3),
978        });
979
980        // Count 1 and 2 should not trigger logging, count 3 should
981        let ex1 = Exchange::new(Message::new("first"));
982        let formatted1 = producer.format_exchange(&ex1, 3);
983        assert!(
984            formatted1.contains("Group of 3:"),
985            "group_size=3 must include count: {formatted1}"
986        );
987        assert!(
988            formatted1.contains("Body: first"),
989            "group log must include body: {formatted1}"
990        );
991    }
992
993    #[test]
994    fn test_group_size_gating_skips_non_multiples() {
995        // group_size=3: counts 1,2,4,5 are skipped; 3 and 6 emit. Exercises the
996        // `call()` early-return path that the format_exchange-only test misses.
997        assert!(!LogProducer::should_emit(1, Some(3)));
998        assert!(!LogProducer::should_emit(2, Some(3)));
999        assert!(LogProducer::should_emit(3, Some(3)));
1000        assert!(!LogProducer::should_emit(4, Some(3)));
1001        assert!(!LogProducer::should_emit(5, Some(3)));
1002        assert!(LogProducer::should_emit(6, Some(3)));
1003    }
1004
1005    #[test]
1006    fn test_group_size_gating_none_emits_all() {
1007        assert!(LogProducer::should_emit(1, None));
1008        assert!(LogProducer::should_emit(100, None));
1009    }
1010
1011    #[test]
1012    fn test_log_config_group_size_param() {
1013        let config = LogConfig::from_uri("log:test?groupSize=10").unwrap();
1014        assert_eq!(config.group_size, Some(10));
1015    }
1016
1017    #[test]
1018    fn test_log_config_group_size_default_unset() {
1019        let config = LogConfig::from_uri("log:test").unwrap();
1020        assert_eq!(config.group_size, None);
1021    }
1022
1023    #[test]
1024    fn test_log_config_show_stream_info_param() {
1025        let config = LogConfig::from_uri("log:test?showStreamInfo=true").unwrap();
1026        assert!(config.show_stream_info);
1027    }
1028
1029    #[test]
1030    fn test_log_config_show_stream_info_default_false() {
1031        let config = LogConfig::from_uri("log:test").unwrap();
1032        assert!(!config.show_stream_info);
1033    }
1034}