1use 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#[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#[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 pub category: String,
81 #[uri_param(
83 kind = "enum:TRACE,DEBUG,INFO,WARN,WARNING,ERROR",
84 default = "INFO",
85 desc = "Log level"
86 )]
87 pub level: LogLevel,
88 #[uri_param(
90 name = "showHeaders",
91 default = "false",
92 desc = "Include headers in log output"
93 )]
94 pub show_headers: bool,
95 #[uri_param(
97 name = "showBody",
98 default = "true",
99 desc = "Include body in log output"
100 )]
101 pub show_body: bool,
102 #[uri_param(
105 name = "maxChars",
106 desc = "Truncate body to N characters. Omit for no limit"
107 )]
108 pub max_chars: Option<usize>,
109 #[uri_param(
113 name = "logMask",
114 default = "false",
115 desc = "Redact sensitive headers and body"
116 )]
117 pub log_mask: bool,
118 #[uri_param(
120 name = "showStreamInfo",
121 default = "false",
122 desc = "Show stream origin info"
123 )]
124 pub show_stream_info: bool,
125 #[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
213pub 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
254struct 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#[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 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 _ => "[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 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 LogLevel::Error => error!("{msg}"),
428 }
429
430 Ok(exchange)
431 })
432 }
433}
434
435impl LogProducer {
436 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#[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 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 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 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 }
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 #[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 #[test]
918 fn test_log_stream_show_info() {
919 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 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 #[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 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 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}