1use async_trait::async_trait;
7use camel_language_api::{Body, Exchange, ExpressionErrorClass, Value};
8use camel_language_api::{Expression, Language, LanguageError, Predicate};
9use jsonpath_rust::parser::model::JpQuery;
10use jsonpath_rust::parser::parse_json_path;
11use jsonpath_rust::query::js_path_process;
12use serde_json::Value as JsonValue;
13#[cfg(test)]
14thread_local! {
15 static COMPILE_COUNT: std::cell::Cell<usize> = const { std::cell::Cell::new(0) };
16}
17
18const DEFAULT_MAX_DEPTH: usize = 64;
20
21const DEFAULT_MAX_INPUT_BYTES: usize = 16 * 1024 * 1024; #[derive(Debug, Clone)]
27pub struct JsonPathConfig {
28 pub max_input_bytes: Option<usize>,
31 pub max_depth: Option<usize>,
34}
35
36impl JsonPathConfig {
37 fn effective_max_depth(&self) -> usize {
39 self.max_depth.unwrap_or(DEFAULT_MAX_DEPTH)
40 }
41}
42
43impl Default for JsonPathConfig {
44 fn default() -> Self {
45 Self {
46 max_input_bytes: Some(DEFAULT_MAX_INPUT_BYTES),
47 max_depth: None,
48 }
49 }
50}
51
52pub struct JsonPathLanguage {
54 config: JsonPathConfig,
55}
56
57impl JsonPathLanguage {
58 pub fn new() -> Self {
60 Self {
61 config: JsonPathConfig::default(),
62 }
63 }
64
65 pub fn with_config(config: JsonPathConfig) -> Self {
67 Self { config }
68 }
69}
70
71impl Default for JsonPathLanguage {
72 fn default() -> Self {
73 Self::new()
74 }
75}
76
77struct JsonPathExpression {
78 query: JpQuery,
79 config: JsonPathConfig,
80}
81
82struct JsonPathPredicate {
83 query: JpQuery,
84 config: JsonPathConfig,
85}
86
87fn check_depth(value: &JsonValue, max_depth: usize) -> Result<(), LanguageError> {
89 fn recurse(value: &JsonValue, max_depth: usize, current: usize) -> Result<(), LanguageError> {
90 if current > max_depth {
91 return Err(LanguageError::EvalError(format!(
92 "JSON nesting depth {current} exceeds limit of {max_depth}"
93 )));
94 }
95 match value {
96 JsonValue::Object(map) => {
97 for v in map.values() {
98 recurse(v, max_depth, current + 1)?;
99 }
100 Ok(())
101 }
102 JsonValue::Array(arr) => {
103 for v in arr {
104 recurse(v, max_depth, current + 1)?;
105 }
106 Ok(())
107 }
108 _ => Ok(()),
109 }
110 }
111 recurse(value, max_depth, 0)
112}
113
114fn extract_json(exchange: &Exchange, config: &JsonPathConfig) -> Result<JsonValue, LanguageError> {
120 let max_depth = config.effective_max_depth();
121
122 match &exchange.input.body {
123 Body::Json(v) => {
124 check_depth(v, max_depth)?;
126 Ok(v.clone())
127 }
128 Body::Text(s) => {
129 if let Some(limit) = config.max_input_bytes
131 && s.len() > limit
132 {
133 return Err(LanguageError::EvalError(format!(
134 "input size {} bytes exceeds limit of {limit} bytes",
135 s.len()
136 )));
137 }
138 let value: JsonValue =
141 serde_json::from_str(s).map_err(|_| LanguageError::EvalFailure {
142 class: ExpressionErrorClass::Conversion,
143 position: None,
144 detail: Some("body is not valid JSON".to_string()),
145 })?;
146 check_depth(&value, max_depth)?;
147 Ok(value)
148 }
149 other => other
150 .clone()
151 .try_into_json()
152 .map_err(|e| {
153 LanguageError::EvalError(format!("body is not JSON and cannot be coerced: {e}"))
154 })
155 .and_then(|b| match b {
156 Body::Json(v) => {
157 check_depth(&v, max_depth)?;
158 Ok(v)
159 }
160 _ => Err(LanguageError::EvalError(
161 "body coercion did not produce JSON".into(),
162 )),
163 }),
164 }
165}
166
167fn run_query(query: &JpQuery, json: &JsonValue) -> Result<JsonValue, LanguageError> {
168 let results = js_path_process(query, json)
169 .map_err(|e| LanguageError::EvalError(format!("jsonpath query '{query}' failed: {e}")))?;
170 let values: Vec<&JsonValue> = results.into_iter().map(|r| r.val).collect();
171 Ok(match values.len() {
172 0 => JsonValue::Null,
173 1 => values[0].clone(),
174 _ => JsonValue::Array(values.into_iter().cloned().collect()),
175 })
176}
177
178#[async_trait]
179impl Expression for JsonPathExpression {
180 async fn evaluate(&self, exchange: &Exchange) -> Result<Value, LanguageError> {
185 let json = extract_json(exchange, &self.config)?;
186 run_query(&self.query, &json)
187 }
188}
189
190#[async_trait]
191impl Predicate for JsonPathPredicate {
192 async fn matches(&self, exchange: &Exchange) -> Result<bool, LanguageError> {
193 let json = extract_json(exchange, &self.config)?;
194 let result = run_query(&self.query, &json)?;
195 match &result {
199 JsonValue::Bool(b) => Ok(*b),
200 other => Err(LanguageError::TypeMismatch {
201 expected: "bool".to_string(),
202 actual: json_type_name(other).to_string(),
203 position: None,
204 }),
205 }
206 }
207}
208
209fn json_type_name(value: &JsonValue) -> &'static str {
212 match value {
213 JsonValue::Null => "null",
214 JsonValue::Bool(_) => "bool",
215 JsonValue::Number(_) => "number",
216 JsonValue::String(_) => "string",
217 JsonValue::Array(_) => "array",
218 JsonValue::Object(_) => "object",
219 }
220}
221
222impl Language for JsonPathLanguage {
223 fn name(&self) -> &'static str {
224 "jsonpath"
225 }
226
227 fn create_expression(&self, script: &str) -> Result<Box<dyn Expression>, LanguageError> {
228 if !script.starts_with('$') {
229 return Err(LanguageError::ParseError {
230 expr: script.to_string(),
231 reason: "JsonPath expression must start with '$'".into(),
232 });
233 }
234 let parsed = parse_json_path(script).map_err(|e| LanguageError::ParseError {
235 expr: script.to_string(),
236 reason: e.to_string(),
237 })?;
238 #[cfg(test)]
239 {
240 COMPILE_COUNT.with(|c| c.set(c.get() + 1));
241 }
242 Ok(Box::new(JsonPathExpression {
243 query: parsed,
244 config: self.config.clone(),
245 }))
246 }
247
248 fn create_predicate(&self, script: &str) -> Result<Box<dyn Predicate>, LanguageError> {
249 if !script.starts_with('$') {
250 return Err(LanguageError::ParseError {
251 expr: script.to_string(),
252 reason: "JsonPath expression must start with '$'".into(),
253 });
254 }
255 let parsed = parse_json_path(script).map_err(|e| LanguageError::ParseError {
256 expr: script.to_string(),
257 reason: e.to_string(),
258 })?;
259 #[cfg(test)]
260 {
261 COMPILE_COUNT.with(|c| c.set(c.get() + 1));
262 }
263 Ok(Box::new(JsonPathPredicate {
264 query: parsed,
265 config: self.config.clone(),
266 }))
267 }
268}
269
270#[cfg(test)]
271mod tests {
272 use super::*;
273 use camel_language_api::Message;
274
275 async fn exchange_with_json(json: &str) -> Exchange {
276 let value: JsonValue = serde_json::from_str(json).unwrap();
277 Exchange::new(Message::new(Body::Json(value)))
278 }
279
280 async fn exchange_with_text_body(text: &str) -> Exchange {
281 Exchange::new(Message::new(Body::Text(text.to_string())))
282 }
283
284 async fn empty_exchange() -> Exchange {
285 Exchange::new(Message::default())
286 }
287
288 async fn default_lang() -> JsonPathLanguage {
289 JsonPathLanguage::new()
290 }
291
292 #[tokio::test]
293 async fn expression_simple_path() {
294 let lang = default_lang().await;
295 let expr = lang.create_expression("$.store.name").unwrap();
296 let ex = exchange_with_json(r#"{"store":{"name":"books"}}"#).await;
297 let result = expr.evaluate(&ex).await.unwrap();
298 assert_eq!(result, JsonValue::String("books".to_string()));
299 }
300
301 #[tokio::test]
302 async fn expression_nested_path() {
303 let lang = default_lang().await;
304 let expr = lang.create_expression("$.a.b.c").unwrap();
305 let ex = exchange_with_json(r#"{"a":{"b":{"c":42}}}"#).await;
306 let result = expr.evaluate(&ex).await.unwrap();
307 assert_eq!(result, JsonValue::Number(42.into()));
308 }
309
310 #[tokio::test]
311 async fn expression_array_index() {
312 let lang = default_lang().await;
313 let expr = lang.create_expression("$.items[0]").unwrap();
314 let ex = exchange_with_json(r#"{"items":["a","b","c"]}"#).await;
315 let result = expr.evaluate(&ex).await.unwrap();
316 assert_eq!(result, JsonValue::String("a".to_string()));
317 }
318
319 #[tokio::test]
320 async fn expression_wildcard() {
321 let lang = default_lang().await;
322 let expr = lang.create_expression("$.items[*].name").unwrap();
323 let ex = exchange_with_json(r#"{"items":[{"name":"a"},{"name":"b"}]}"#).await;
324 let result = expr.evaluate(&ex).await.unwrap();
325 assert_eq!(
326 result,
327 JsonValue::Array(vec![
328 JsonValue::String("a".to_string()),
329 JsonValue::String("b".to_string())
330 ])
331 );
332 }
333
334 #[tokio::test]
335 async fn expression_root_path() {
336 let lang = default_lang().await;
337 let expr = lang.create_expression("$").unwrap();
338 let ex = exchange_with_json(r#"{"x":1}"#).await;
339 let result = expr.evaluate(&ex).await.unwrap();
340 assert_eq!(result["x"], JsonValue::Number(1.into()));
341 }
342
343 #[tokio::test]
344 async fn expression_text_body_with_valid_json() {
345 let lang = default_lang().await;
346 let expr = lang.create_expression("$.name").unwrap();
347 let ex = exchange_with_text_body(r#"{"name":"test"}"#).await;
348 let result = expr.evaluate(&ex).await.unwrap();
349 assert_eq!(result, JsonValue::String("test".to_string()));
350 }
351
352 #[tokio::test]
353 async fn expression_empty_body_is_error() {
354 let lang = default_lang().await;
355 let expr = lang.create_expression("$.x").unwrap();
356 let ex = empty_exchange().await;
357 let result = expr.evaluate(&ex).await;
358 assert!(result.is_err());
359 }
360
361 #[tokio::test]
362 async fn expression_invalid_jsonpath_syntax() {
363 let lang = default_lang().await;
364 let result = lang.create_expression("$[invalid");
365 let err = match result {
366 Err(e) => e,
367 Ok(_) => panic!("expected ParseError"),
368 };
369 match err {
370 LanguageError::ParseError { expr, reason } => {
371 assert!(!expr.is_empty());
372 assert!(!reason.is_empty());
373 }
374 other => panic!("expected ParseError, got {other:?}"),
375 }
376 }
377
378 #[tokio::test]
381 async fn expression_without_dollar_prefix_is_rejected() {
382 let lang = default_lang().await;
383 let result = lang.create_expression("store.name");
384 assert!(result.is_err(), "expected error for missing $ prefix");
385 let err = match result {
386 Err(e) => e,
387 Ok(_) => panic!("expected ParseError"),
388 };
389 match err {
390 LanguageError::ParseError { expr, reason } => {
391 assert_eq!(expr, "store.name");
392 assert!(
393 reason.contains("'$'"),
394 "reason should mention '$', got: {reason}"
395 );
396 }
397 other => panic!("expected ParseError, got {other:?}"),
398 }
399 }
400
401 #[tokio::test]
402 async fn predicate_without_dollar_prefix_is_rejected() {
403 let lang = default_lang().await;
404 let result = lang.create_predicate("store.name");
405 assert!(result.is_err(), "expected error for missing $ prefix");
406 let err = match result {
407 Err(e) => e,
408 Ok(_) => panic!("expected ParseError"),
409 };
410 match err {
411 LanguageError::ParseError { reason, .. } => {
412 assert!(
413 reason.contains("'$'"),
414 "reason should mention '$', got: {reason}"
415 );
416 }
417 other => panic!("expected ParseError, got {other:?}"),
418 }
419 }
420
421 #[tokio::test]
424 async fn expression_deeply_nested_path() {
425 let lang = default_lang().await;
426 let expr = lang.create_expression("$.a.b.c.d").unwrap();
427 let ex = exchange_with_json(r#"{"a":{"b":{"c":{"d":"deep"}}}}"#).await;
428 let result = expr.evaluate(&ex).await.unwrap();
429 assert_eq!(result, JsonValue::String("deep".to_string()));
430 }
431
432 #[tokio::test]
433 async fn expression_array_index_nested() {
434 let lang = default_lang().await;
435 let expr = lang.create_expression("$.data.items[1].name").unwrap();
436 let ex = exchange_with_json(
437 r#"{"data":{"items":[{"name":"first"},{"name":"second"},{"name":"third"}]}}"#,
438 )
439 .await;
440 let result = expr.evaluate(&ex).await.unwrap();
441 assert_eq!(result, JsonValue::String("second".to_string()));
442 }
443
444 #[tokio::test]
445 async fn jsonpath_predicate_non_bool_is_type_mismatch() {
446 let lang = default_lang().await;
450 for (json_body, actual) in [
451 (r#"{"val":"false"}"#, "string"),
452 (r#"{"val":"x"}"#, "string"),
453 (r#"{"val":0}"#, "number"),
454 (r#"{"val":1}"#, "number"),
455 (r#"{"val":[]}"#, "array"),
456 (r#"{"val":{}}"#, "object"),
457 (r#"{"other":1}"#, "null"),
458 ] {
459 let pred = lang.create_predicate("$.val").unwrap();
460 let ex = exchange_with_json(json_body).await;
461 match pred.matches(&ex).await {
462 Err(LanguageError::TypeMismatch {
463 expected,
464 actual: got,
465 position: None,
466 }) => {
467 assert_eq!(expected, "bool");
468 assert_eq!(got, actual, "body {json_body}");
469 }
470 other => panic!("expected TypeMismatch for {json_body}, got: {other:?}"),
471 }
472 }
473
474 let pred = lang.create_predicate("$.active").unwrap();
476 assert!(
477 pred.matches(&exchange_with_json(r#"{"active":true}"#).await)
478 .await
479 .unwrap()
480 );
481 let pred = lang.create_predicate("$.active").unwrap();
482 assert!(
483 !pred
484 .matches(&exchange_with_json(r#"{"active":false}"#).await)
485 .await
486 .unwrap()
487 );
488 }
489
490 #[tokio::test]
491 async fn jsonpath_invalid_body_error_is_conversion_class_no_snippet() {
492 let lang = default_lang().await;
493 let expr = lang.create_expression("$.key").unwrap();
494 let mut ex = Exchange::new(Message::default());
495 ex.input.body = Body::Text("SECRETVAL not json".to_string());
496 let err = expr.evaluate(&ex).await.expect_err("must fail");
497 match &err {
498 LanguageError::EvalFailure {
499 class: ExpressionErrorClass::Conversion,
500 position: None,
501 detail: Some(detail),
502 } => {
503 assert_eq!(detail, "body is not valid JSON");
504 }
505 other => panic!("expected Conversion EvalFailure, got: {other:?}"),
506 }
507 assert!(!err.to_string().contains("SECRETVAL"), "body leaked: {err}");
508 }
509
510 #[tokio::test]
513 async fn oversized_input_is_rejected() {
514 let config = JsonPathConfig {
515 max_input_bytes: Some(100),
516 ..Default::default()
517 };
518 let lang = JsonPathLanguage::with_config(config);
519 let expr = lang.create_expression("$.key").unwrap();
520 let big_value = "x".repeat(200);
522 let big_json = format!(r#"{{"key":"{}"}}"#, big_value);
523 assert!(big_json.len() > 100);
524 let ex = exchange_with_text_body(&big_json).await;
525 let result = expr.evaluate(&ex).await;
526 assert!(
527 result.is_err(),
528 "expected error for oversized input, got {result:?}"
529 );
530 }
531
532 #[tokio::test]
533 async fn input_under_limit_is_accepted() {
534 let config = JsonPathConfig {
535 max_input_bytes: Some(1024),
536 ..Default::default()
537 };
538 let lang = JsonPathLanguage::with_config(config);
539 let expr = lang.create_expression("$.key").unwrap();
540 let ex = exchange_with_text_body(r#"{"key":"value"}"#).await;
541 let result = expr.evaluate(&ex).await;
542 assert!(
543 result.is_ok(),
544 "expected success for input under limit, got {result:?}"
545 );
546 }
547
548 #[tokio::test]
549 async fn deeply_nested_input_is_rejected() {
550 let config = JsonPathConfig {
551 max_depth: Some(5),
552 ..Default::default()
553 };
554 let lang = JsonPathLanguage::with_config(config);
555 let expr = lang.create_expression("$.a").unwrap();
556 let mut json = "1".to_string();
558 for _ in 0..10 {
559 json = format!(r#"{{"a":{json}}}"#);
560 }
561 let ex = exchange_with_text_body(&json).await;
562 let result = expr.evaluate(&ex).await;
563 assert!(
564 result.is_err(),
565 "expected error for deeply nested input, got {result:?}"
566 );
567 }
568
569 #[tokio::test]
570 async fn nesting_within_depth_limit_is_accepted() {
571 let config = JsonPathConfig {
572 max_depth: Some(10),
573 ..Default::default()
574 };
575 let lang = JsonPathLanguage::with_config(config);
576 let expr = lang.create_expression("$.a").unwrap();
577 let mut json = "1".to_string();
579 for _ in 0..5 {
580 json = format!(r#"{{"a":{json}}}"#);
581 }
582 let ex = exchange_with_text_body(&json).await;
583 let result = expr.evaluate(&ex).await;
584 assert!(
585 result.is_ok(),
586 "expected success for nesting within limit, got {result:?}"
587 );
588 }
589
590 #[tokio::test]
591 async fn default_config_has_safe_defaults() {
592 let config = JsonPathConfig::default();
594 assert_eq!(config.max_input_bytes, Some(16 * 1024 * 1024));
595 assert_eq!(config.max_depth, None);
597 }
598
599 #[tokio::test]
600 async fn default_config_rejects_oversized_text_body() {
601 let lang = JsonPathLanguage::new(); let expr = lang.create_expression("$.key").unwrap();
604 let big_value = "x".repeat(16 * 1024 * 1024 + 10);
605 let big_json = format!(r#"{{"key":"{}"}}"#, big_value);
606 let ex = exchange_with_text_body(&big_json).await;
607 let result = expr.evaluate(&ex).await;
608 assert!(
609 result.is_err(),
610 "default config must reject oversized input, got {result:?}"
611 );
612 }
613
614 #[tokio::test]
615 async fn oversized_input_also_rejected_for_predicate() {
616 let config = JsonPathConfig {
617 max_input_bytes: Some(100),
618 ..Default::default()
619 };
620 let lang = JsonPathLanguage::with_config(config);
621 let pred = lang.create_predicate("$.key").unwrap();
622 let big_value = "x".repeat(200);
623 let big_json = format!(r#"{{"key":"{}"}}"#, big_value);
624 let ex = exchange_with_text_body(&big_json).await;
625 let result = pred.matches(&ex).await;
626 assert!(
627 result.is_err(),
628 "expected error for oversized input in predicate, got {result:?}"
629 );
630 }
631
632 #[tokio::test]
633 async fn deeply_nested_input_rejected_for_predicate() {
634 let config = JsonPathConfig {
635 max_depth: Some(3),
636 ..Default::default()
637 };
638 let lang = JsonPathLanguage::with_config(config);
639 let pred = lang.create_predicate("$.a").unwrap();
640 let mut json = "1".to_string();
641 for _ in 0..5 {
642 json = format!(r#"{{"a":{json}}}"#);
643 }
644 let ex = exchange_with_text_body(&json).await;
645 let result = pred.matches(&ex).await;
646 assert!(
647 result.is_err(),
648 "expected error for deeply nested input in predicate, got {result:?}"
649 );
650 }
651
652 #[tokio::test]
653 async fn body_json_no_input_size_check_but_depth_checked() {
654 let config = JsonPathConfig {
657 max_input_bytes: Some(10), max_depth: Some(3),
659 };
660 let lang = JsonPathLanguage::with_config(config);
661 let expr = lang.create_expression("$.a").unwrap();
662 let mut json_str = "1".to_string();
664 for _ in 0..5 {
665 json_str = format!(r#"{{"a":{json_str}}}"#);
666 }
667 let ex = exchange_with_json(&json_str).await;
668 let result = expr.evaluate(&ex).await;
669 assert!(
670 result.is_err(),
671 "expected depth error for pre-parsed JSON, got {result:?}"
672 );
673 }
674
675 #[test]
681 fn compile_once_types_are_send_sync() {
682 fn assert_send_sync<T: Send + Sync>() {}
683 assert_send_sync::<JsonPathExpression>();
684 assert_send_sync::<JsonPathPredicate>();
685 assert_send_sync::<JpQuery>();
686 }
687
688 #[tokio::test]
697 async fn compilation_happens_only_once_for_expression() {
698 let before = COMPILE_COUNT.with(std::cell::Cell::get);
699 let lang = JsonPathLanguage::new();
700 let expr = lang.create_expression("$.foo.bar").unwrap();
701 let after_create = COMPILE_COUNT.with(std::cell::Cell::get);
702 assert_eq!(
703 after_create,
704 before + 1,
705 "create_expression must compile exactly once"
706 );
707
708 for _ in 0..50 {
710 let ex = exchange_with_json(r#"{"foo":{"bar":"baz"}}"#).await;
711 let result = expr.evaluate(&ex).await.unwrap();
712 assert_eq!(result, JsonValue::String("baz".to_string()));
713 }
714 assert_eq!(
715 COMPILE_COUNT.with(std::cell::Cell::get),
716 after_create,
717 "evaluate must NOT trigger re-compilation"
718 );
719 }
720
721 #[tokio::test]
722 async fn compilation_happens_only_once_for_predicate() {
723 let before = COMPILE_COUNT.with(std::cell::Cell::get);
724 let lang = JsonPathLanguage::new();
725 let pred = lang.create_predicate("$.active").unwrap();
726 let after_create = COMPILE_COUNT.with(std::cell::Cell::get);
727 assert_eq!(
728 after_create,
729 before + 1,
730 "create_predicate must compile exactly once"
731 );
732
733 for _ in 0..50 {
734 let ex = exchange_with_json(r#"{"active":true}"#).await;
735 assert!(pred.matches(&ex).await.unwrap());
736 }
737 assert_eq!(
738 COMPILE_COUNT.with(std::cell::Cell::get),
739 after_create,
740 "matches must NOT trigger re-compilation"
741 );
742 }
743}