1use crate::event::SpanEvent;
12use crate::ingest::IngestSource;
13
14pub const MAX_JSON_DEPTH: usize = 32;
23
24#[must_use]
38pub fn exceeds_max_depth(raw: &[u8]) -> bool {
39 let mut depth: usize = 0;
40 let mut in_string = false;
41 let mut escape = false;
42 for &b in raw {
43 if in_string {
44 advance_string_state(b, &mut in_string, &mut escape);
45 continue;
46 }
47 if bump_depth(b, &mut depth, &mut in_string) {
48 return true;
49 }
50 }
51 false
52}
53
54#[inline]
59fn advance_string_state(b: u8, in_string: &mut bool, escape: &mut bool) {
60 if *escape {
61 *escape = false;
62 } else if b == b'\\' {
63 *escape = true;
64 } else if b == b'"' {
65 *in_string = false;
66 }
67}
68
69#[inline]
73fn bump_depth(b: u8, depth: &mut usize, in_string: &mut bool) -> bool {
74 match b {
75 b'"' => *in_string = true,
76 b'[' | b'{' => {
77 *depth += 1;
78 if *depth > MAX_JSON_DEPTH {
79 return true;
80 }
81 }
82 b']' | b'}' => *depth = depth.saturating_sub(1),
83 _ => {}
84 }
85 false
86}
87
88#[derive(Debug, Clone, Copy, PartialEq, Eq)]
93#[non_exhaustive]
94pub enum InputFormat {
95 Native,
97 Otlp,
99 Jaeger,
101 Zipkin,
103}
104
105pub struct JsonIngest {
107 max_size: usize,
108}
109
110impl JsonIngest {
111 #[must_use]
112 pub const fn new(max_size: usize) -> Self {
113 Self { max_size }
114 }
115}
116
117impl IngestSource for JsonIngest {
118 type Error = JsonIngestError;
119
120 fn ingest(&self, raw: &[u8]) -> Result<Vec<SpanEvent>, Self::Error> {
121 if raw.len() > self.max_size {
122 return Err(JsonIngestError::PayloadTooLarge {
123 size: raw.len(),
124 max: self.max_size,
125 });
126 }
127
128 if exceeds_max_depth(raw) {
133 return Err(JsonIngestError::PayloadTooDeep {
134 max_depth: MAX_JSON_DEPTH,
135 });
136 }
137
138 match detect_format(raw) {
139 InputFormat::Otlp => {
140 type OtlpRequest =
150 opentelemetry_proto::tonic::collector::trace::v1::ExportTraceServiceRequest;
151 let mut events = Vec::new();
152 let mut parsed_any = false;
153 let mut offset = 0;
154 while offset < raw.len() {
155 let mut stream = serde_json::Deserializer::from_slice(&raw[offset..])
156 .into_iter::<OtlpRequest>();
157 match stream.next() {
158 None => break,
159 Some(Ok(request)) => {
160 parsed_any = true;
161 events.extend(crate::ingest::otlp::convert_otlp_request(&request));
162 offset += stream.byte_offset();
163 }
164 Some(Err(e)) if is_missing_values(&e) => {
170 let mut retry = serde_json::Deserializer::from_slice(&raw[offset..])
171 .into_iter::<serde_json::Value>();
172 let Some(Ok(mut value)) = retry.next() else {
173 return Err(JsonIngestError::Parse(e));
174 };
175 normalize_otlp_json(&mut value);
176 let request: OtlpRequest =
177 serde_json::from_value(value).map_err(JsonIngestError::Parse)?;
178 parsed_any = true;
179 events.extend(crate::ingest::otlp::convert_otlp_request(&request));
180 offset += retry.byte_offset();
181 }
182 Some(Err(e)) if e.is_eof() && parsed_any => {
188 tracing::warn!(
189 "ignoring truncated trailing OTLP JSON document \
190 (live or rotated file-exporter dump?)"
191 );
192 break;
193 }
194 Some(Err(e)) => return Err(JsonIngestError::Parse(e)),
195 }
196 }
197 Ok(events)
198 }
199 InputFormat::Jaeger => {
200 let ingest = crate::ingest::jaeger::JaegerIngest::new(self.max_size);
201 ingest
202 .ingest(raw)
203 .map_err(|e| JsonIngestError::Format(e.to_string()))
204 }
205 InputFormat::Zipkin => {
206 let ingest = crate::ingest::zipkin::ZipkinIngest::new(self.max_size);
207 ingest
208 .ingest(raw)
209 .map_err(|e| JsonIngestError::Format(e.to_string()))
210 }
211 InputFormat::Native => {
212 let mut events: Vec<SpanEvent> =
213 serde_json::from_slice(raw).map_err(JsonIngestError::Parse)?;
214 for event in &mut events {
219 if let Some(region) = event.cloud_region.as_deref()
220 && !crate::score::carbon::is_valid_region_id(region)
221 {
222 event.cloud_region = None;
223 }
224 crate::event::sanitize_span_event(event);
225 }
226 Ok(events)
227 }
228 }
229 }
230}
231
232fn is_missing_values(e: &serde_json::Error) -> bool {
238 e.classify() == serde_json::error::Category::Data
239 && e.to_string().contains("missing field `values`")
240}
241
242fn normalize_otlp_json(value: &mut serde_json::Value) {
248 match value {
249 serde_json::Value::Object(map) => {
250 for key in ["arrayValue", "kvlistValue"] {
251 if let Some(serde_json::Value::Object(inner)) = map.get_mut(key)
252 && !inner.contains_key("values")
253 {
254 inner.insert("values".to_string(), serde_json::Value::Array(Vec::new()));
255 }
256 }
257 for v in map.values_mut() {
258 normalize_otlp_json(v);
259 }
260 }
261 serde_json::Value::Array(items) => items.iter_mut().for_each(normalize_otlp_json),
262 _ => {}
263 }
264}
265
266#[must_use]
271pub fn detect_format(raw: &[u8]) -> InputFormat {
272 let peek = std::str::from_utf8(&raw[..raw.len().min(1024)]).unwrap_or("");
273
274 if peek.trim_start().starts_with('{') {
282 let mut saw_data_key = false;
283 for key in TopLevelKeys::new(peek) {
284 match key {
285 "resourceSpans" | "resource_spans" => return InputFormat::Otlp,
290 "data" => saw_data_key = true,
293 _ => {}
294 }
295 }
296 if saw_data_key {
297 let deeper = std::str::from_utf8(&raw[..raw.len().min(4096)]).unwrap_or("");
298 if deeper.contains("\"spans\"") {
299 return InputFormat::Jaeger;
300 }
301 }
302 }
303
304 if peek.trim_start().starts_with('[')
306 && peek.contains("\"traceId\"")
307 && peek.contains("\"localEndpoint\"")
308 {
309 return InputFormat::Zipkin;
310 }
311
312 InputFormat::Native
313}
314
315struct TopLevelKeys<'a> {
322 bytes: &'a [u8],
323 pos: usize,
324 depth: usize,
325}
326
327impl<'a> TopLevelKeys<'a> {
328 fn new(peek: &'a str) -> Self {
329 Self {
330 bytes: peek.as_bytes(),
331 pos: 0,
332 depth: 0,
333 }
334 }
335
336 fn scan_string(&mut self) -> Option<(usize, usize)> {
341 let start = self.pos + 1;
342 let mut i = start;
343 let mut escape = false;
344 while i < self.bytes.len() {
345 let c = self.bytes[i];
346 if escape {
347 escape = false;
348 } else if c == b'\\' {
349 escape = true;
350 } else if c == b'"' {
351 break;
352 }
353 i += 1;
354 }
355 if i >= self.bytes.len() {
356 return None; }
358 self.pos = i + 1;
359 Some((start, i))
360 }
361
362 fn colon_follows(&self) -> bool {
365 let mut j = self.pos;
366 while j < self.bytes.len() && self.bytes[j].is_ascii_whitespace() {
367 j += 1;
368 }
369 j < self.bytes.len() && self.bytes[j] == b':'
370 }
371}
372
373impl<'a> Iterator for TopLevelKeys<'a> {
374 type Item = &'a str;
375
376 fn next(&mut self) -> Option<&'a str> {
377 while self.pos < self.bytes.len() {
378 match self.bytes[self.pos] {
379 b'{' | b'[' => {
380 self.depth += 1;
381 self.pos += 1;
382 }
383 b'}' | b']' => {
384 self.depth = self.depth.saturating_sub(1);
385 self.pos += 1;
386 }
387 b'"' => {
388 let (start, end) = self.scan_string()?;
389 if self.depth == 1
390 && self.colon_follows()
391 && let Ok(key) = std::str::from_utf8(&self.bytes[start..end])
392 {
393 return Some(key);
394 }
395 }
396 _ => self.pos += 1,
397 }
398 }
399 None
400 }
401}
402
403#[derive(Debug, thiserror::Error)]
407#[non_exhaustive]
408pub enum JsonIngestError {
409 #[error("payload too large: {size} bytes exceeds maximum of {max} bytes")]
410 PayloadTooLarge { size: usize, max: usize },
411 #[error(
412 "payload nesting exceeds maximum depth of {max_depth} (defense against deeply-nested attacker payloads)"
413 )]
414 PayloadTooDeep { max_depth: usize },
415 #[error("JSON parse error: {0}")]
416 Parse(#[from] serde_json::Error),
417 #[error("format detection error: {0}")]
418 Format(String),
419}
420
421#[cfg(test)]
422mod tests {
423 use super::*;
424 use core::assert_matches;
425
426 #[test]
427 fn rejects_oversized_payload() {
428 let ingest = JsonIngest::new(10);
429 let result = ingest.ingest(&[0u8; 100]);
430 assert!(result.is_err());
431 }
432
433 #[test]
434 fn parses_empty_array() {
435 let ingest = JsonIngest::new(1_048_576);
436 let events = ingest.ingest(b"[]").unwrap();
437 assert!(events.is_empty());
438 }
439
440 #[test]
441 fn detect_native_format() {
442 let json = r#"[{"type": "sql", "target": "SELECT 1"}]"#;
443 assert_eq!(detect_format(json.as_bytes()), InputFormat::Native);
444 }
445
446 #[test]
447 fn detect_jaeger_format() {
448 let json = r#"{"data": [{"traceID": "abc", "spans": [], "processes": {}}]}"#;
449 assert_eq!(detect_format(json.as_bytes()), InputFormat::Jaeger);
450 }
451
452 #[test]
453 fn detect_zipkin_format() {
454 let json = r#"[{"traceId": "abc", "id": "s1", "localEndpoint": {"serviceName": "svc"}}]"#;
455 assert_eq!(detect_format(json.as_bytes()), InputFormat::Zipkin);
456 }
457
458 #[test]
459 fn detect_empty_array_is_native() {
460 assert_eq!(detect_format(b"[]"), InputFormat::Native);
461 }
462
463 #[test]
464 fn detect_invalid_json_falls_to_native() {
465 assert_eq!(detect_format(b"not json"), InputFormat::Native);
466 }
467
468 #[test]
469 fn auto_ingest_jaeger() {
470 let json = r#"{
471 "data": [{
472 "traceID": "t1",
473 "spans": [{
474 "spanID": "s1",
475 "operationName": "op",
476 "references": [],
477 "startTime": 1720621921123000,
478 "duration": 500,
479 "processID": "p1",
480 "tags": [
481 {"key": "db.statement", "value": "SELECT 1"},
482 {"key": "db.system", "value": "pg"}
483 ]
484 }],
485 "processes": {"p1": {"serviceName": "svc"}}
486 }]
487 }"#;
488 let ingest = JsonIngest::new(1_048_576);
489 let events = ingest.ingest(json.as_bytes()).unwrap();
490 assert_eq!(events.len(), 1);
491 assert_eq!(events[0].target, "SELECT 1");
492 }
493
494 #[test]
495 fn auto_ingest_zipkin() {
496 let json = r#"[{
497 "traceId": "t1",
498 "id": "s1",
499 "name": "query",
500 "timestamp": 1720621921123000,
501 "duration": 500,
502 "localEndpoint": {"serviceName": "svc"},
503 "tags": {"db.statement": "SELECT 1", "db.system": "pg"}
504 }]"#;
505 let ingest = JsonIngest::new(1_048_576);
506 let events = ingest.ingest(json.as_bytes()).unwrap();
507 assert_eq!(events.len(), 1);
508 assert_eq!(events[0].target, "SELECT 1");
509 }
510
511 fn otlp_request_json(trace_id: &str, statement: &str) -> String {
515 otlp_request_json_with_attrs(trace_id, statement, "")
516 }
517
518 fn otlp_request_json_with_attrs(trace_id: &str, statement: &str, extra_attrs: &str) -> String {
522 format!(
523 r#"{{"resourceSpans":[{{"resource":{{"attributes":[{{"key":"service.name","value":{{"stringValue":"svc"}}}}]}},"scopeSpans":[{{"spans":[{{"traceId":"{trace_id}","spanId":"eee19b7ec3c1b174","name":"db-query","kind":3,"startTimeUnixNano":"1720621921000000000","endTimeUnixNano":"1720621921000500000","attributes":[{{"key":"db.statement","value":{{"stringValue":"{statement}"}}}},{{"key":"db.system","value":{{"stringValue":"postgresql"}}}}{extra_attrs}]}}]}}]}}]}}"#
524 )
525 }
526
527 #[test]
528 fn detect_otlp_format() {
529 let json = r#"{"resourceSpans": [{"scopeSpans": []}]}"#;
530 assert_eq!(detect_format(json.as_bytes()), InputFormat::Otlp);
531 }
532
533 #[test]
534 fn detect_jaeger_wins_over_stray_resource_spans_literal() {
535 let json = r#"{
539 "data": [{
540 "traceID": "t1",
541 "spans": [{
542 "spanID": "s1",
543 "operationName": "export resourceSpans",
544 "references": [],
545 "startTime": 1720621921123000,
546 "duration": 500,
547 "processID": "p1",
548 "tags": [{"key": "note", "value": "handles \"resourceSpans\" batches"}]
549 }],
550 "processes": {"p1": {"serviceName": "collector"}}
551 }]
552 }"#;
553 assert_eq!(detect_format(json.as_bytes()), InputFormat::Jaeger);
554 }
555
556 #[test]
557 fn detect_otlp_wins_over_stray_data_literal() {
558 let json = r#"{"resourceSpans":[{"resource":{"attributes":[{"key":"data","value":{"stringValue":"data"}}]},"scopeSpans":[{"spans":[]}]}]}"#;
563 assert_eq!(detect_format(json.as_bytes()), InputFormat::Otlp);
564 }
565
566 #[test]
567 fn top_level_keys_ignores_nested_keys_and_string_values() {
568 let json = r#"{"a": {"nested": 1}, "b": ["data", {"c": 2}], "d": "resourceSpans"}"#;
569 let keys: Vec<&str> = TopLevelKeys::new(json).collect();
570 assert_eq!(keys, ["a", "b", "d"]);
571 }
572
573 #[test]
574 fn detect_otlp_snake_case_routes_to_otlp() {
575 let json = r#"{"resource_spans": []}"#;
579 assert_eq!(detect_format(json.as_bytes()), InputFormat::Otlp);
580 let ingest = JsonIngest::new(1_048_576);
581 assert_matches!(
582 ingest.ingest(json.as_bytes()),
583 Err(JsonIngestError::Parse(_))
584 );
585 }
586
587 #[test]
588 fn auto_ingest_otlp() {
589 let json = otlp_request_json("5b8efff798038103d269b633813fc60c", "SELECT 1");
590 let ingest = JsonIngest::new(1_048_576);
591 let events = ingest.ingest(json.as_bytes()).unwrap();
592 assert_eq!(events.len(), 1);
593 assert_eq!(events[0].target, "SELECT 1");
594 assert_eq!(events[0].service.as_ref(), "svc");
595 assert_eq!(events[0].trace_id, "5b8efff798038103d269b633813fc60c");
596 }
597
598 #[test]
599 fn otlp_empty_array_value_ingests() {
600 let json = otlp_request_json_with_attrs(
603 "5b8efff798038103d269b633813fc60c",
604 "SELECT 1",
605 r#",{"key":"tags","value":{"arrayValue":{}}}"#,
606 );
607 let ingest = JsonIngest::new(1_048_576);
608 let events = ingest.ingest(json.as_bytes()).unwrap();
609 assert_eq!(events.len(), 1);
610 assert_eq!(events[0].target, "SELECT 1");
611 }
612
613 #[test]
614 fn otlp_empty_kvlist_value_ingests() {
615 let json = otlp_request_json_with_attrs(
617 "5b8efff798038103d269b633813fc60c",
618 "SELECT 1",
619 r#",{"key":"meta","value":{"kvlistValue":{}}}"#,
620 );
621 let ingest = JsonIngest::new(1_048_576);
622 let events = ingest.ingest(json.as_bytes()).unwrap();
623 assert_eq!(events.len(), 1);
624 assert_eq!(events[0].target, "SELECT 1");
625 }
626
627 #[test]
628 fn normalize_fills_missing_array_values() {
629 use opentelemetry_proto::tonic::common::v1::{AnyValue, any_value};
634 let mut value: serde_json::Value = serde_json::from_str(r#"{"arrayValue":{}}"#).unwrap();
635 normalize_otlp_json(&mut value);
636 let any: AnyValue = serde_json::from_value(value).unwrap();
637 let Some(any_value::Value::ArrayValue(av)) = any.value else {
638 panic!("expected ArrayValue variant");
639 };
640 assert!(av.values.is_empty());
641 }
642
643 #[test]
644 fn otlp_issue_81_repro_line_ingests() {
645 let json = r#"{"resourceSpans":[{"resource":{"attributes":[{"key":"service.name","value":{"stringValue":"svc-a"}}]},"scopeSpans":[{"scope":{"name":"repro"},"spans":[{"traceId":"5b8efff798038103d269b633813fc60c","spanId":"eee19b7ec3c1b174","name":"GET /x","kind":2,"startTimeUnixNano":"1783678644000000000","endTimeUnixNano":"1783678644100000000","attributes":[{"key":"empty.list","value":{"arrayValue":{}}}]}]}]}]}"#;
649 let ingest = JsonIngest::new(1_048_576);
650 assert!(ingest.ingest(json.as_bytes()).is_ok());
651 }
652
653 #[test]
654 fn otlp_empty_array_value_mid_ndjson_continues() {
655 let line1 = otlp_request_json_with_attrs(
658 "0af7651916cd43dd8448eb211c80319c",
659 "SELECT 1",
660 r#",{"key":"tags","value":{"arrayValue":{}}}"#,
661 );
662 let line2 = otlp_request_json("1bf7651916cd43dd8448eb211c80319d", "SELECT 2");
663 let json = format!("{line1}\n{line2}\n");
664 let ingest = JsonIngest::new(1_048_576);
665 let events = ingest.ingest(json.as_bytes()).unwrap();
666 assert_eq!(events.len(), 2);
667 assert_eq!(events[0].target, "SELECT 1");
668 assert_eq!(events[1].target, "SELECT 2");
669 }
670
671 #[test]
672 fn otlp_type_wrong_truncated_tail_still_fails() {
673 let full = otlp_request_json("0af7651916cd43dd8448eb211c80319c", "SELECT 1");
677 let json = format!("{full}\n{{\"resourceSpans\":123");
678 let ingest = JsonIngest::new(1_048_576);
679 assert_matches!(
680 ingest.ingest(json.as_bytes()),
681 Err(JsonIngestError::Parse(_))
682 );
683 }
684
685 #[test]
686 fn auto_ingest_otlp_ndjson() {
687 let json = format!(
689 "{}\n{}\n",
690 otlp_request_json("0af7651916cd43dd8448eb211c80319c", "SELECT 1"),
691 otlp_request_json("1bf7651916cd43dd8448eb211c80319d", "SELECT 2"),
692 );
693 let ingest = JsonIngest::new(1_048_576);
694 let events = ingest.ingest(json.as_bytes()).unwrap();
695 assert_eq!(events.len(), 2);
696 assert_eq!(events[0].target, "SELECT 1");
697 assert_eq!(events[1].target, "SELECT 2");
698 }
699
700 #[test]
701 fn auto_ingest_otlp_ndjson_tolerates_truncated_tail() {
702 let full = otlp_request_json("0af7651916cd43dd8448eb211c80319c", "SELECT 1");
706 let truncated = &full[..full.len() / 2];
707 let json = format!("{full}\n{truncated}");
708 let ingest = JsonIngest::new(1_048_576);
709 let events = ingest.ingest(json.as_bytes()).unwrap();
710 assert_eq!(events.len(), 1);
711 assert_eq!(events[0].target, "SELECT 1");
712 }
713
714 #[test]
715 fn auto_ingest_otlp_truncated_only_payload_still_fails() {
716 let full = otlp_request_json("0af7651916cd43dd8448eb211c80319c", "SELECT 1");
719 let truncated = &full[..full.len() / 2];
720 let ingest = JsonIngest::new(1_048_576);
721 assert_matches!(
722 ingest.ingest(truncated.as_bytes()),
723 Err(JsonIngestError::Parse(_))
724 );
725 }
726
727 #[test]
728 fn auto_ingest_otlp_mid_stream_garbage_still_fails() {
729 let full = otlp_request_json("0af7651916cd43dd8448eb211c80319c", "SELECT 1");
732 let json = format!("{full}\n{{\"resourceSpans\": 42}}\n{full}");
733 let ingest = JsonIngest::new(1_048_576);
734 assert_matches!(
735 ingest.ingest(json.as_bytes()),
736 Err(JsonIngestError::Parse(_))
737 );
738 }
739
740 #[test]
741 fn auto_ingest_otlp_pretty_printed_single_object() {
742 let compact = otlp_request_json("5b8efff798038103d269b633813fc60c", "SELECT 1");
745 let value: serde_json::Value = serde_json::from_str(&compact).unwrap();
746 let pretty = serde_json::to_string_pretty(&value).unwrap();
747 let ingest = JsonIngest::new(1_048_576);
748 let events = ingest.ingest(pretty.as_bytes()).unwrap();
749 assert_eq!(events.len(), 1);
750 assert_eq!(events[0].target, "SELECT 1");
751 }
752
753 #[test]
754 fn deeply_nested_otlp_payload_is_rejected() {
755 let depth = MAX_JSON_DEPTH + 4;
758 let mut payload = String::from(
759 r#"{"resourceSpans":[{"scopeSpans":[{"spans":[{"attributes":[{"key":"a","value":{"arrayValue":{"values":["#,
760 );
761 for _ in 0..depth {
762 payload.push('[');
763 }
764 for _ in 0..depth {
765 payload.push(']');
766 }
767 payload.push_str("]}}}]}]}]}]}");
768 let ingest = JsonIngest::new(1_048_576);
769 let result = ingest.ingest(payload.as_bytes());
770 assert!(
771 matches!(result, Err(JsonIngestError::PayloadTooDeep { .. })),
772 "deeply-nested OTLP input must be rejected: {result:?}"
773 );
774 }
775
776 fn native_event_with_cloud_region(cloud_region: &str) -> String {
779 format!(
780 r#"[{{
781 "timestamp": "2025-07-10T14:32:01.123Z",
782 "trace_id": "trace-1",
783 "span_id": "span-1",
784 "service": "order-svc",
785 "cloud_region": {cr},
786 "type": "sql",
787 "operation": "SELECT",
788 "target": "SELECT 1",
789 "duration_us": 1000,
790 "source": {{
791 "endpoint": "POST /api/orders/42/submit",
792 "method": "OrderService::create_order"
793 }}
794 }}]"#,
795 cr = serde_json::to_string(cloud_region).unwrap()
796 )
797 }
798
799 #[test]
800 fn native_json_valid_cloud_region_preserved() {
801 let json = native_event_with_cloud_region("eu-west-3");
803 let ingest = JsonIngest::new(1_048_576);
804 let events = ingest.ingest(json.as_bytes()).unwrap();
805 assert_eq!(events.len(), 1);
806 assert_eq!(events[0].cloud_region.as_deref(), Some("eu-west-3"));
807 }
808
809 #[test]
810 fn native_json_invalid_cloud_region_is_sanitized_to_none() {
811 let json = native_event_with_cloud_region("eu-west-3\n2026 WARN fake alert");
815 let ingest = JsonIngest::new(1_048_576);
816 let events = ingest.ingest(json.as_bytes()).unwrap();
817 assert_eq!(events.len(), 1);
818 assert!(
819 events[0].cloud_region.is_none(),
820 "cloud_region with control char must be sanitized"
821 );
822 }
823
824 #[test]
825 fn native_json_oversized_cloud_region_sanitized() {
826 let long_region = "a".repeat(65);
828 let json = native_event_with_cloud_region(&long_region);
829 let ingest = JsonIngest::new(1_048_576);
830 let events = ingest.ingest(json.as_bytes()).unwrap();
831 assert!(events[0].cloud_region.is_none());
832 }
833
834 #[test]
835 fn native_json_cloud_region_with_space_sanitized() {
836 let json = native_event_with_cloud_region("eu west 3");
837 let ingest = JsonIngest::new(1_048_576);
838 let events = ingest.ingest(json.as_bytes()).unwrap();
839 assert!(events[0].cloud_region.is_none());
840 }
841
842 #[test]
843 fn native_json_cloud_region_with_dot_sanitized() {
844 let json = native_event_with_cloud_region("eu.west.3");
846 let ingest = JsonIngest::new(1_048_576);
847 let events = ingest.ingest(json.as_bytes()).unwrap();
848 assert!(events[0].cloud_region.is_none());
849 }
850
851 #[test]
852 fn deeply_nested_native_payload_is_rejected_below_stack_overflow() {
853 let depth = MAX_JSON_DEPTH + 4;
856 let mut payload = String::with_capacity(depth * 2);
857 for _ in 0..depth {
858 payload.push('[');
859 }
860 for _ in 0..depth {
861 payload.push(']');
862 }
863 let ingest = JsonIngest::new(1_048_576);
864 let result = ingest.ingest(payload.as_bytes());
865 assert_matches!(result, Err(JsonIngestError::PayloadTooDeep { .. }));
866 }
867
868 #[test]
869 fn deeply_nested_jaeger_payload_is_rejected() {
870 let depth = MAX_JSON_DEPTH + 4;
874 let mut payload = String::from(r#"{"data":[{"spans":[{"tags":["#);
875 for _ in 0..depth {
876 payload.push('[');
877 }
878 for _ in 0..depth {
879 payload.push(']');
880 }
881 payload.push_str("]}]}]}");
882 let ingest = JsonIngest::new(1_048_576);
883 let result = ingest.ingest(payload.as_bytes());
884 assert!(
885 matches!(result, Err(JsonIngestError::PayloadTooDeep { .. })),
886 "deeply-nested Jaeger input must be rejected: {result:?}"
887 );
888 }
889
890 #[test]
891 fn deeply_nested_zipkin_payload_is_rejected() {
892 let depth = MAX_JSON_DEPTH + 4;
894 let mut payload = String::from(
895 r#"[{"traceId":"abc","localEndpoint":{"serviceName":"s"},"annotations":["#,
896 );
897 for _ in 0..depth {
898 payload.push('[');
899 }
900 for _ in 0..depth {
901 payload.push(']');
902 }
903 payload.push_str("]}]");
904 let ingest = JsonIngest::new(1_048_576);
905 let result = ingest.ingest(payload.as_bytes());
906 assert!(
907 matches!(result, Err(JsonIngestError::PayloadTooDeep { .. })),
908 "deeply-nested Zipkin input must be rejected: {result:?}"
909 );
910 }
911
912 #[test]
919 fn native_ingest_accepts_input_at_depth_31() {
920 let mut payload = String::with_capacity(64);
922 for _ in 0..31 {
923 payload.push('[');
924 }
925 for _ in 0..31 {
926 payload.push(']');
927 }
928 let ingest = JsonIngest::new(1_048_576);
929 let result = ingest.ingest(payload.as_bytes());
930 assert!(
931 !matches!(result, Err(JsonIngestError::PayloadTooDeep { .. })),
932 "depth 31 must not be rejected by the depth guard, got: {result:?}"
933 );
934 }
935
936 #[test]
937 fn native_ingest_rejects_input_at_depth_33() {
938 let mut payload = String::with_capacity(68);
939 for _ in 0..33 {
940 payload.push('[');
941 }
942 for _ in 0..33 {
943 payload.push(']');
944 }
945 let ingest = JsonIngest::new(1_048_576);
946 assert_matches!(
947 ingest.ingest(payload.as_bytes()),
948 Err(JsonIngestError::PayloadTooDeep { .. })
949 );
950 }
951
952 #[test]
953 fn jaeger_ingest_accepts_input_at_depth_31() {
954 let inner = 25;
957 let mut payload = String::from(r#"{"data":[{"spans":[{"tags":["#);
958 for _ in 0..inner {
959 payload.push('[');
960 }
961 for _ in 0..inner {
962 payload.push(']');
963 }
964 payload.push_str("]}]}]}");
965 let ingest = JsonIngest::new(1_048_576);
966 let result = ingest.ingest(payload.as_bytes());
967 assert!(
968 !matches!(result, Err(JsonIngestError::PayloadTooDeep { .. })),
969 "Jaeger depth 31 must not be rejected by the depth guard, got: {result:?}"
970 );
971 }
972
973 #[test]
974 fn jaeger_ingest_rejects_input_at_depth_33() {
975 let inner = 27;
977 let mut payload = String::from(r#"{"data":[{"spans":[{"tags":["#);
978 for _ in 0..inner {
979 payload.push('[');
980 }
981 for _ in 0..inner {
982 payload.push(']');
983 }
984 payload.push_str("]}]}]}");
985 let ingest = JsonIngest::new(1_048_576);
986 assert_matches!(
987 ingest.ingest(payload.as_bytes()),
988 Err(JsonIngestError::PayloadTooDeep { .. })
989 );
990 }
991
992 #[test]
993 fn zipkin_ingest_accepts_input_at_depth_31() {
994 let inner = 28;
997 let mut payload = String::from(
998 r#"[{"traceId":"abc","localEndpoint":{"serviceName":"s"},"annotations":["#,
999 );
1000 for _ in 0..inner {
1001 payload.push('[');
1002 }
1003 for _ in 0..inner {
1004 payload.push(']');
1005 }
1006 payload.push_str("]}]");
1007 let ingest = JsonIngest::new(1_048_576);
1008 let result = ingest.ingest(payload.as_bytes());
1009 assert!(
1010 !matches!(result, Err(JsonIngestError::PayloadTooDeep { .. })),
1011 "Zipkin depth 31 must not be rejected by the depth guard, got: {result:?}"
1012 );
1013 }
1014
1015 #[test]
1016 fn zipkin_ingest_rejects_input_at_depth_33() {
1017 let inner = 30;
1019 let mut payload = String::from(
1020 r#"[{"traceId":"abc","localEndpoint":{"serviceName":"s"},"annotations":["#,
1021 );
1022 for _ in 0..inner {
1023 payload.push('[');
1024 }
1025 for _ in 0..inner {
1026 payload.push(']');
1027 }
1028 payload.push_str("]}]");
1029 let ingest = JsonIngest::new(1_048_576);
1030 assert_matches!(
1031 ingest.ingest(payload.as_bytes()),
1032 Err(JsonIngestError::PayloadTooDeep { .. })
1033 );
1034 }
1035
1036 #[test]
1037 fn depth_scan_ignores_brackets_inside_strings() {
1038 let json = native_event_with_cloud_region("eu-west-3").replace(
1043 "\"SELECT 1\"",
1044 "\"SELECT * FROM t WHERE col = '[[[[[[[[[[[[[[[[[[[[[[[[[[[[[[[[[[[[[[[[[[[[]'\"",
1045 );
1046 let ingest = JsonIngest::new(1_048_576);
1047 let events = ingest
1048 .ingest(json.as_bytes())
1049 .expect("string-internal brackets must not trigger the depth guard");
1050 assert_eq!(events.len(), 1);
1051 }
1052}