1use std::collections::{BTreeMap, HashMap};
16use std::fmt;
17use std::str::FromStr;
18use std::sync::Arc;
19
20use axum::Extension;
21use axum::extract::{Path, Query, State};
22use axum::http::{HeaderMap, StatusCode as HttpStatusCode};
23use axum::response::IntoResponse;
24use axum_extra::TypedHeader;
25use common_catalog::consts::{PARENT_SPAN_ID_COLUMN, TRACE_TABLE_NAME};
26use common_error::ext::ErrorExt;
27use common_error::status_code::StatusCode;
28use common_query::{Output, OutputData};
29use common_recordbatch::util;
30use common_telemetry::{debug, error, tracing, warn};
31use headers::UserAgent;
32use serde::{Deserialize, Deserializer, Serialize, de};
33use serde_json::Value as JsonValue;
34use session::context::{Channel, QueryContext};
35use snafu::{OptionExt, ResultExt};
36
37use crate::error::{
38 CollectRecordbatchSnafu, Error, InvalidJaegerQuerySnafu, Result, status_code_to_http_status,
39};
40use crate::http::HttpRecordsOutput;
41use crate::http::extractor::TraceTableName;
42use crate::metrics::METRIC_JAEGER_QUERY_ELAPSED;
43use crate::otlp::trace::{
44 DURATION_NANO_COLUMN, KEY_OTEL_SCOPE_NAME, KEY_OTEL_SCOPE_VERSION, KEY_OTEL_STATUS_CODE,
45 KEY_OTEL_STATUS_ERROR_KEY, KEY_OTEL_STATUS_MESSAGE, KEY_OTEL_TRACE_STATE, KEY_SERVICE_NAME,
46 KEY_SPAN_KIND, RESOURCE_ATTRIBUTES_COLUMN, SCOPE_NAME_COLUMN, SCOPE_VERSION_COLUMN,
47 SERVICE_NAME_COLUMN, SPAN_ATTRIBUTES_COLUMN, SPAN_EVENTS_COLUMN, SPAN_ID_COLUMN,
48 SPAN_KIND_COLUMN, SPAN_KIND_PREFIX, SPAN_LINKS_COLUMN, SPAN_NAME_COLUMN, SPAN_STATUS_CODE,
49 SPAN_STATUS_ERROR, SPAN_STATUS_MESSAGE_COLUMN, SPAN_STATUS_PREFIX, SPAN_STATUS_UNSET,
50 TIMESTAMP_COLUMN, TRACE_ID_COLUMN, TRACE_STATE_COLUMN,
51};
52use crate::query_handler::JaegerQueryHandlerRef;
53
54pub const JAEGER_QUERY_TABLE_NAME_KEY: &str = "jaeger_query_table_name";
55
56const REF_TYPE_CHILD_OF: &str = "CHILD_OF";
57const SPAN_KIND_TIME_FMTS: [&str; 2] = ["%Y-%m-%d %H:%M:%S%.6f%z", "%Y-%m-%d %H:%M:%S%.9f%z"];
58
59const TRACE_NOT_FOUND_ERROR_CODE: i32 = 404;
60const TRACE_NOT_FOUND_ERROR_MSG: &str = "trace not found";
61
62#[derive(Default, Debug, Serialize, Deserialize, PartialEq)]
65pub struct JaegerAPIResponse {
66 pub data: Option<JaegerData>,
67 pub total: usize,
68 pub limit: usize,
69 pub offset: usize,
70 pub errors: Vec<JaegerAPIError>,
71}
72
73impl JaegerAPIResponse {
74 pub fn trace_not_found() -> Self {
75 Self {
76 data: None,
77 total: 0,
78 limit: 0,
79 offset: 0,
80 errors: vec![JaegerAPIError {
81 code: TRACE_NOT_FOUND_ERROR_CODE,
82 msg: TRACE_NOT_FOUND_ERROR_MSG.to_string(),
83 trace_id: None,
84 }],
85 }
86 }
87}
88
89#[derive(Debug, Serialize, Deserialize, PartialEq)]
91#[serde(untagged)]
92pub enum JaegerData {
93 ServiceNames(Vec<String>),
94 OperationsNames(Vec<String>),
95 Operations(Vec<Operation>),
96 Traces(Vec<Trace>),
97}
98
99#[derive(Default, Debug, Serialize, Deserialize, PartialEq)]
101#[serde(rename_all = "camelCase")]
102pub struct JaegerAPIError {
103 pub code: i32,
104 pub msg: String,
105 #[serde(skip_serializing_if = "Option::is_none")]
106 pub trace_id: Option<String>,
107}
108
109#[derive(Debug, Default, Serialize, Deserialize, PartialEq)]
111#[serde(rename_all = "camelCase")]
112pub struct Operation {
113 pub name: String,
114 #[serde(skip_serializing_if = "Option::is_none")]
115 pub span_kind: Option<String>,
116}
117
118#[derive(Debug, Default, Serialize, Deserialize, PartialEq)]
120#[serde(rename_all = "camelCase")]
121pub struct Trace {
122 #[serde(rename = "traceID")]
123 pub trace_id: String,
124 pub spans: Vec<Span>,
125
126 #[serde(skip_serializing_if = "HashMap::is_empty")]
127 pub processes: HashMap<String, Process>,
128
129 #[serde(skip_serializing_if = "Vec::is_empty")]
130 pub warnings: Vec<String>,
131}
132
133#[derive(Debug, Default, Serialize, Deserialize, PartialEq)]
135#[serde(rename_all = "camelCase")]
136pub struct Span {
137 #[serde(rename = "traceID")]
138 pub trace_id: String,
139
140 #[serde(rename = "spanID")]
141 pub span_id: String,
142
143 #[serde(rename = "parentSpanID")]
144 #[serde(skip_serializing_if = "String::is_empty")]
145 pub parent_span_id: String,
146
147 #[serde(skip_serializing_if = "Option::is_none")]
148 pub flags: Option<u32>,
149
150 pub operation_name: String,
151 pub references: Vec<Reference>,
152 pub start_time: u64, pub duration: u64, pub tags: Vec<KeyValue>,
155 pub logs: Vec<Log>,
156
157 #[serde(rename = "processID")]
158 #[serde(skip_serializing_if = "String::is_empty")]
159 pub process_id: String,
160
161 #[serde(skip_serializing_if = "Option::is_none")]
162 pub process: Option<Process>,
163
164 #[serde(skip_serializing_if = "Vec::is_empty")]
165 pub warnings: Vec<String>,
166}
167
168#[derive(Debug, Serialize, Deserialize, PartialEq)]
170#[serde(rename_all = "camelCase")]
171pub struct Reference {
172 #[serde(rename = "traceID")]
173 pub trace_id: String,
174 #[serde(rename = "spanID")]
175 pub span_id: String,
176 pub ref_type: String,
177}
178
179#[derive(Debug, Serialize, Deserialize, PartialEq)]
181#[serde(rename_all = "camelCase")]
182pub struct Process {
183 pub service_name: String,
184 pub tags: Vec<KeyValue>,
185}
186
187#[derive(Debug, Serialize, Deserialize, PartialEq)]
189#[serde(rename_all = "camelCase")]
190pub struct Log {
191 pub timestamp: u64,
192 pub fields: Vec<KeyValue>,
193}
194
195#[derive(Debug, Serialize, Deserialize, PartialEq)]
197#[serde(rename_all = "camelCase")]
198pub struct KeyValue {
199 pub key: String,
200 #[serde(rename = "type")]
201 pub value_type: ValueType,
202 pub value: Value,
203}
204
205#[derive(Debug, Serialize, Deserialize, PartialEq)]
207#[serde(untagged)]
208#[serde(rename_all = "camelCase")]
209pub enum Value {
210 String(String),
211 Int64(i64),
212 Float64(f64),
213 Boolean(bool),
214 Binary(Vec<u8>),
215}
216
217#[derive(Debug, Serialize, Deserialize, PartialEq)]
219#[serde(rename_all = "lowercase")]
220pub enum ValueType {
221 String,
222 Int64,
223 Float64,
224 Boolean,
225 Binary,
226}
227
228#[derive(Default, Debug, Serialize, Deserialize)]
230#[serde(rename_all = "camelCase")]
231pub struct JaegerQueryParams {
232 #[serde(rename = "service")]
234 pub service_name: Option<String>,
235
236 #[serde(rename = "operation")]
238 pub operation_name: Option<String>,
239
240 #[serde(default, deserialize_with = "empty_string_as_none")]
242 pub limit: Option<usize>,
243
244 pub start: Option<i64>,
246
247 pub end: Option<i64>,
249
250 #[serde(default, deserialize_with = "empty_string_as_none")]
252 pub max_duration: Option<String>,
253
254 #[serde(default, deserialize_with = "empty_string_as_none")]
256 pub min_duration: Option<String>,
257
258 pub tags: Option<String>,
262
263 pub span_kind: Option<String>,
265}
266
267fn empty_string_as_none<'de, D, T>(de: D) -> Result<Option<T>, D::Error>
269where
270 D: Deserializer<'de>,
271 T: FromStr,
272 T::Err: fmt::Display,
273{
274 let opt = Option::<String>::deserialize(de)?;
275 match opt.as_deref() {
276 None | Some("") => Ok(None),
277 Some(s) => FromStr::from_str(s).map_err(de::Error::custom).map(Some),
278 }
279}
280
281fn update_query_context(query_ctx: &mut QueryContext, table_name: Option<String>) {
282 query_ctx.set_channel(Channel::Jaeger);
284 if let Some(table) = table_name {
285 query_ctx.set_extension(JAEGER_QUERY_TABLE_NAME_KEY, table);
286 }
287}
288
289impl QueryTraceParams {
290 fn from_jaeger_query_params(query_params: JaegerQueryParams) -> Result<Self> {
291 let mut internal_query_params: QueryTraceParams = QueryTraceParams {
292 service_name: query_params.service_name.context(InvalidJaegerQuerySnafu {
293 reason: "service_name is required".to_string(),
294 })?,
295 operation_name: query_params.operation_name,
296 start_time: query_params.start.map(|start| start * 1000),
298 end_time: query_params.end.map(|end| end * 1000),
299 ..Default::default()
300 };
301
302 if let Some(max_duration) = query_params.max_duration {
303 let duration = humantime::parse_duration(&max_duration).map_err(|e| {
304 InvalidJaegerQuerySnafu {
305 reason: format!("parse maxDuration '{}' failed: {}", max_duration, e),
306 }
307 .build()
308 })?;
309 internal_query_params.max_duration = Some(duration.as_nanos() as u64);
310 }
311
312 if let Some(min_duration) = query_params.min_duration {
313 let duration = humantime::parse_duration(&min_duration).map_err(|e| {
314 InvalidJaegerQuerySnafu {
315 reason: format!("parse minDuration '{}' failed: {}", min_duration, e),
316 }
317 .build()
318 })?;
319 internal_query_params.min_duration = Some(duration.as_nanos() as u64);
320 }
321
322 if let Some(tags) = query_params.tags {
323 let mut tags_map: HashMap<String, JsonValue> =
325 serde_json::from_str(&tags).map_err(|e| {
326 InvalidJaegerQuerySnafu {
327 reason: format!("parse tags '{}' failed: {}", tags, e),
328 }
329 .build()
330 })?;
331 for (_, v) in tags_map.iter_mut() {
332 if let Some(number) = convert_string_to_number(v) {
333 *v = number;
334 }
335 if let Some(boolean) = convert_string_to_boolean(v) {
336 *v = boolean;
337 }
338 }
339 internal_query_params.tags = Some(tags_map);
340 }
341
342 internal_query_params.limit = query_params.limit;
343
344 Ok(internal_query_params)
345 }
346}
347
348#[derive(Debug, Default, PartialEq)]
349pub struct QueryTraceParams {
350 pub service_name: String,
351 pub operation_name: Option<String>,
352
353 pub limit: Option<usize>,
355
356 pub tags: Option<HashMap<String, JsonValue>>,
358
359 pub start_time: Option<i64>,
361 pub end_time: Option<i64>,
362 pub min_duration: Option<u64>,
363 pub max_duration: Option<u64>,
364
365 pub user_agent: TraceUserAgent,
367}
368
369#[derive(Debug, Default, PartialEq, Eq)]
370pub enum TraceUserAgent {
371 Grafana,
372 #[default]
375 Jaeger,
376}
377
378impl From<UserAgent> for TraceUserAgent {
379 fn from(value: UserAgent) -> Self {
380 let ua_str = value.as_str().to_lowercase();
381 debug!("received user agent: {}", ua_str);
382 if ua_str.contains("grafana") {
383 Self::Grafana
384 } else {
385 Self::Jaeger
386 }
387 }
388}
389
390#[axum_macros::debug_handler]
392#[tracing::instrument(skip_all, fields(protocol = "jaeger", request_type = "get_services"))]
393pub async fn handle_get_services(
394 State(handler): State<JaegerQueryHandlerRef>,
395 Query(query_params): Query<JaegerQueryParams>,
396 Extension(mut query_ctx): Extension<QueryContext>,
397 TraceTableName(table_name): TraceTableName,
398) -> impl IntoResponse {
399 debug!(
400 "Received Jaeger '/api/services' request, query_params: {:?}, query_ctx: {:?}",
401 query_params, query_ctx
402 );
403
404 query_ctx.set_channel(Channel::Jaeger);
405 if let Some(table) = table_name {
406 query_ctx.set_extension(JAEGER_QUERY_TABLE_NAME_KEY, table);
407 }
408
409 let query_ctx = Arc::new(query_ctx);
410 let db = query_ctx.get_db_string();
411
412 let _timer = METRIC_JAEGER_QUERY_ELAPSED
414 .with_label_values(&[db.as_str(), "/api/services"])
415 .start_timer();
416
417 match handler.get_services(query_ctx).await {
418 Ok(output) => match covert_to_records(output).await {
419 Ok(Some(records)) => match services_from_records(records) {
420 Ok(services) => {
421 let services_num = services.len();
422 (
423 HttpStatusCode::OK,
424 axum::Json(JaegerAPIResponse {
425 data: Some(JaegerData::ServiceNames(services)),
426 total: services_num,
427 ..Default::default()
428 }),
429 )
430 }
431 Err(err) => {
432 error!("Failed to get services: {:?}", err);
433 error_response(err)
434 }
435 },
436 Ok(None) => (HttpStatusCode::OK, axum::Json(JaegerAPIResponse::default())),
437 Err(err) => {
438 error!("Failed to get services: {:?}", err);
439 error_response(err)
440 }
441 },
442 Err(err) => handle_query_error(err, "Failed to get services", &db),
443 }
444}
445
446#[axum_macros::debug_handler]
448#[tracing::instrument(skip_all, fields(protocol = "jaeger", request_type = "get_trace"))]
449pub async fn handle_get_trace(
450 State(handler): State<JaegerQueryHandlerRef>,
451 Path(trace_id): Path<String>,
452 Query(query_params): Query<JaegerQueryParams>,
453 Extension(mut query_ctx): Extension<QueryContext>,
454 TraceTableName(table_name): TraceTableName,
455) -> impl IntoResponse {
456 debug!(
457 "Received Jaeger '/api/traces/{}' request, query_params: {:?}, query_ctx: {:?}",
458 trace_id, query_params, query_ctx
459 );
460
461 update_query_context(&mut query_ctx, table_name);
462 let query_ctx = Arc::new(query_ctx);
463 let db = query_ctx.get_db_string();
464
465 let _timer = METRIC_JAEGER_QUERY_ELAPSED
467 .with_label_values(&[db.as_str(), "/api/traces"])
468 .start_timer();
469
470 let start_time_ns = query_params.start.map(|start_us| start_us * 1000);
472 let end_time_ns = query_params.end.map(|end_us| end_us * 1000);
473
474 let output = match handler
475 .get_trace(
476 query_ctx,
477 &trace_id,
478 start_time_ns,
479 end_time_ns,
480 query_params.limit,
481 )
482 .await
483 {
484 Ok(output) => output,
485 Err(err) => {
486 return handle_query_error(
487 err,
488 &format!("Failed to get trace for '{}'", trace_id),
489 &db,
490 );
491 }
492 };
493
494 match covert_to_records(output).await {
495 Ok(Some(records)) => match traces_from_records(records) {
496 Ok(traces) if traces.is_empty() => (
497 HttpStatusCode::NOT_FOUND,
498 axum::Json(JaegerAPIResponse::trace_not_found()),
499 ),
500 Ok(traces) => (
501 HttpStatusCode::OK,
502 axum::Json(JaegerAPIResponse {
503 data: Some(JaegerData::Traces(traces)),
504 ..Default::default()
505 }),
506 ),
507 Err(err) => {
508 error!("Failed to get trace '{}': {:?}", trace_id, err);
509 error_response(err)
510 }
511 },
512 Ok(None) => (
513 HttpStatusCode::NOT_FOUND,
514 axum::Json(JaegerAPIResponse::trace_not_found()),
515 ),
516 Err(err) => {
517 error!("Failed to get trace '{}': {:?}", trace_id, err);
518 error_response(err)
519 }
520 }
521}
522
523#[axum_macros::debug_handler]
525#[tracing::instrument(skip_all, fields(protocol = "jaeger", request_type = "find_traces"))]
526pub async fn handle_find_traces(
527 State(handler): State<JaegerQueryHandlerRef>,
528 Query(query_params): Query<JaegerQueryParams>,
529 Extension(mut query_ctx): Extension<QueryContext>,
530 TraceTableName(table_name): TraceTableName,
531 optional_user_agent: Option<TypedHeader<UserAgent>>,
532) -> impl IntoResponse {
533 debug!(
534 "Received Jaeger '/api/traces' request, query_params: {:?}, query_ctx: {:?}",
535 query_params, query_ctx
536 );
537
538 update_query_context(&mut query_ctx, table_name);
539 let query_ctx = Arc::new(query_ctx);
540 let db = query_ctx.get_db_string();
541
542 let _timer = METRIC_JAEGER_QUERY_ELAPSED
544 .with_label_values(&[db.as_str(), "/api/traces"])
545 .start_timer();
546
547 match QueryTraceParams::from_jaeger_query_params(query_params) {
548 Ok(mut query_params) => {
549 if let Some(TypedHeader(user_agent)) = optional_user_agent {
550 query_params.user_agent = user_agent.into();
551 }
552 let output = handler.find_traces(query_ctx, query_params).await;
553 match output {
554 Ok(output) => match covert_to_records(output).await {
555 Ok(Some(records)) => match traces_from_records(records) {
556 Ok(traces) => (
557 HttpStatusCode::OK,
558 axum::Json(JaegerAPIResponse {
559 data: Some(JaegerData::Traces(traces)),
560 ..Default::default()
561 }),
562 ),
563 Err(err) => {
564 error!("Failed to find traces: {:?}", err);
565 error_response(err)
566 }
567 },
568 Ok(None) => (HttpStatusCode::OK, axum::Json(JaegerAPIResponse::default())),
569 Err(err) => error_response(err),
570 },
571 Err(err) => handle_query_error(err, "Failed to find traces", &db),
572 }
573 }
574 Err(e) => error_response(e),
575 }
576}
577
578#[axum_macros::debug_handler]
580#[tracing::instrument(skip_all, fields(protocol = "jaeger", request_type = "get_operations"))]
581pub async fn handle_get_operations(
582 State(handler): State<JaegerQueryHandlerRef>,
583 Query(query_params): Query<JaegerQueryParams>,
584 Extension(mut query_ctx): Extension<QueryContext>,
585 TraceTableName(table_name): TraceTableName,
586 headers: HeaderMap,
587) -> impl IntoResponse {
588 debug!(
589 "Received Jaeger '/api/operations' request, query_params: {:?}, query_ctx: {:?}, headers: {:?}",
590 query_params, query_ctx, headers
591 );
592
593 if let Some(service_name) = &query_params.service_name {
594 update_query_context(&mut query_ctx, table_name);
595 let query_ctx = Arc::new(query_ctx);
596 let db = query_ctx.get_db_string();
597
598 let _timer = METRIC_JAEGER_QUERY_ELAPSED
600 .with_label_values(&[db.as_str(), "/api/operations"])
601 .start_timer();
602
603 match handler
604 .get_operations(query_ctx, service_name, query_params.span_kind.as_deref())
605 .await
606 {
607 Ok(output) => match covert_to_records(output).await {
608 Ok(Some(records)) => match operations_from_records(records, true) {
609 Ok(operations) => {
610 let total = operations.len();
611 (
612 HttpStatusCode::OK,
613 axum::Json(JaegerAPIResponse {
614 data: Some(JaegerData::Operations(operations)),
615 total,
616 ..Default::default()
617 }),
618 )
619 }
620 Err(err) => {
621 error!("Failed to get operations: {:?}", err);
622 error_response(err)
623 }
624 },
625 Ok(None) => (HttpStatusCode::OK, axum::Json(JaegerAPIResponse::default())),
626 Err(err) => error_response(err),
627 },
628 Err(err) => handle_query_error(
629 err,
630 &format!("Failed to get operations for service '{}'", service_name),
631 &db,
632 ),
633 }
634 } else {
635 (
636 HttpStatusCode::BAD_REQUEST,
637 axum::Json(JaegerAPIResponse {
638 errors: vec![JaegerAPIError {
639 code: 400,
640 msg: "parameter 'service' is required".to_string(),
641 trace_id: None,
642 }],
643 ..Default::default()
644 }),
645 )
646 }
647}
648
649#[axum_macros::debug_handler]
651#[tracing::instrument(
652 skip_all,
653 fields(protocol = "jaeger", request_type = "get_operations_by_service")
654)]
655pub async fn handle_get_operations_by_service(
656 State(handler): State<JaegerQueryHandlerRef>,
657 Path(service_name): Path<String>,
658 Query(query_params): Query<JaegerQueryParams>,
659 Extension(mut query_ctx): Extension<QueryContext>,
660 TraceTableName(table_name): TraceTableName,
661 headers: HeaderMap,
662) -> impl IntoResponse {
663 debug!(
664 "Received Jaeger '/api/services/{}/operations' request, query_params: {:?}, query_ctx: {:?}, headers: {:?}",
665 service_name, query_params, query_ctx, headers
666 );
667
668 update_query_context(&mut query_ctx, table_name);
669 let query_ctx = Arc::new(query_ctx);
670 let db = query_ctx.get_db_string();
671
672 let _timer = METRIC_JAEGER_QUERY_ELAPSED
674 .with_label_values(&[db.as_str(), "/api/services"])
675 .start_timer();
676
677 match handler.get_operations(query_ctx, &service_name, None).await {
678 Ok(output) => match covert_to_records(output).await {
679 Ok(Some(records)) => match operations_from_records(records, false) {
680 Ok(operations) => {
681 let operations: Vec<String> =
682 operations.into_iter().map(|op| op.name).collect();
683 let total = operations.len();
684 (
685 HttpStatusCode::OK,
686 axum::Json(JaegerAPIResponse {
687 data: Some(JaegerData::OperationsNames(operations)),
688 total,
689 ..Default::default()
690 }),
691 )
692 }
693 Err(err) => {
694 error!(
695 "Failed to get operations for service '{}': {:?}",
696 service_name, err
697 );
698 error_response(err)
699 }
700 },
701 Ok(None) => (HttpStatusCode::OK, axum::Json(JaegerAPIResponse::default())),
702 Err(err) => error_response(err),
703 },
704 Err(err) => handle_query_error(
705 err,
706 &format!("Failed to get operations for service '{}'", service_name),
707 &db,
708 ),
709 }
710}
711
712async fn covert_to_records(output: Output) -> Result<Option<HttpRecordsOutput>> {
713 match output.data {
714 OutputData::Stream(stream) => {
715 let records = HttpRecordsOutput::try_new(
716 stream.schema().clone(),
717 util::collect(stream)
718 .await
719 .context(CollectRecordbatchSnafu)?,
720 )?;
721 debug!(
722 "The query records: {}",
723 serde_json::to_string(&records).unwrap()
724 );
725 Ok(Some(records))
726 }
727 _ => Ok(None),
729 }
730}
731
732fn handle_query_error(
733 err: Error,
734 prompt: &str,
735 db: &str,
736) -> (HttpStatusCode, axum::Json<JaegerAPIResponse>) {
737 if err.status_code() == StatusCode::TableNotFound {
739 warn!(
740 "No trace table '{}' found in database '{}'",
741 TRACE_TABLE_NAME, db
742 );
743 (HttpStatusCode::OK, axum::Json(JaegerAPIResponse::default()))
744 } else {
745 error!("{}: {:?}", prompt, err);
746 error_response(err)
747 }
748}
749
750fn error_response(err: Error) -> (HttpStatusCode, axum::Json<JaegerAPIResponse>) {
751 (
752 status_code_to_http_status(&err.status_code()),
753 axum::Json(JaegerAPIResponse {
754 errors: vec![JaegerAPIError {
755 code: err.status_code() as i32,
756 msg: err.to_string(),
757 ..Default::default()
758 }],
759 ..Default::default()
760 }),
761 )
762}
763
764fn traces_from_records(records: HttpRecordsOutput) -> Result<Vec<Trace>> {
765 let mut trace_id_to_processes: HashMap<String, HashMap<String, String>> = HashMap::new();
767 let mut trace_id_to_spans: BTreeMap<String, Vec<Span>> = BTreeMap::new();
770 let mut service_to_resource_attributes: HashMap<String, Vec<KeyValue>> = HashMap::new();
772
773 let is_span_attributes_flatten = !records
774 .schema
775 .column_schemas
776 .iter()
777 .any(|c| c.name == SPAN_ATTRIBUTES_COLUMN);
778
779 for row in records.rows.into_iter() {
780 let mut span = Span::default();
781 let mut service_name = None;
782 let mut parent_span_id = None;
783 let mut resource_tags = vec![];
784
785 for (idx, cell) in row.into_iter().enumerate() {
786 let column_name = &records.schema.column_schemas[idx].name;
788
789 match column_name.as_str() {
790 TRACE_ID_COLUMN => {
791 if let JsonValue::String(trace_id) = cell {
792 span.trace_id = trace_id.clone();
793 trace_id_to_processes.entry(trace_id).or_default();
794 }
795 }
796 TIMESTAMP_COLUMN => {
797 span.start_time = cell.as_u64().context(InvalidJaegerQuerySnafu {
798 reason: "Failed to convert timestamp to u64".to_string(),
799 })? / 1000;
800 }
801 DURATION_NANO_COLUMN => {
802 span.duration = cell.as_u64().context(InvalidJaegerQuerySnafu {
803 reason: "Failed to convert duration to u64".to_string(),
804 })? / 1000;
805 }
806 SERVICE_NAME_COLUMN => {
807 if let JsonValue::String(name) = cell {
808 service_name = Some(name);
809 }
810 }
811 SPAN_NAME_COLUMN => {
812 if let JsonValue::String(span_name) = cell {
813 span.operation_name = span_name;
814 }
815 }
816 SPAN_ID_COLUMN => {
817 if let JsonValue::String(span_id) = cell {
818 span.span_id = span_id;
819 }
820 }
821 SPAN_ATTRIBUTES_COLUMN => {
822 if let JsonValue::Object(span_attrs) = cell {
825 span.tags.extend(object_to_tags(span_attrs));
826 }
827 }
828 RESOURCE_ATTRIBUTES_COLUMN => {
829 if let JsonValue::Object(mut resource_attrs) = cell {
833 resource_attrs.remove(KEY_SERVICE_NAME);
834 resource_tags = object_to_tags(resource_attrs);
835 }
836 }
837 PARENT_SPAN_ID_COLUMN => {
838 if let JsonValue::String(id) = cell
839 && !id.is_empty()
840 {
841 parent_span_id = Some(id);
842 }
843 }
844 SPAN_LINKS_COLUMN => {
845 if let JsonValue::Array(links) = cell {
846 for link in links {
847 if let (Some(trace_id), Some(span_id)) = (
848 link.get("trace_id").and_then(JsonValue::as_str),
849 link.get("span_id").and_then(JsonValue::as_str),
850 ) {
851 span.references.push(Reference {
852 trace_id: trace_id.to_string(),
853 span_id: span_id.to_string(),
854 ref_type: "FOLLOWS_FROM".to_string(),
855 });
856 }
857 }
858 }
859 }
860 SPAN_EVENTS_COLUMN => {
861 if let JsonValue::Array(events) = cell {
862 for event in events {
863 if let JsonValue::Object(mut obj) = event {
864 let Some(action) = obj.get("name").and_then(|v| v.as_str()) else {
865 continue;
866 };
867
868 let Some(t) =
869 obj.get("time").and_then(|t| t.as_str()).and_then(|s| {
870 SPAN_KIND_TIME_FMTS
871 .iter()
872 .find_map(|fmt| {
873 chrono::DateTime::parse_from_str(s, fmt).ok()
874 })
875 .map(|dt| dt.timestamp_micros() as u64)
876 })
877 else {
878 continue;
879 };
880
881 let mut fields = vec![KeyValue {
882 key: "event".to_string(),
883 value_type: ValueType::String,
884 value: Value::String(action.to_string()),
885 }];
886
887 if let Some(JsonValue::Object(attrs)) = obj.remove("attributes") {
889 fields.extend(object_to_tags(attrs));
890 }
891
892 span.logs.push(Log {
893 timestamp: t,
894 fields,
895 });
896 }
897 }
898 }
899 }
900 SCOPE_NAME_COLUMN => {
901 if let JsonValue::String(scope_name) = cell
902 && !scope_name.is_empty()
903 {
904 span.tags.push(KeyValue {
905 key: KEY_OTEL_SCOPE_NAME.to_string(),
906 value_type: ValueType::String,
907 value: Value::String(scope_name),
908 });
909 }
910 }
911 SCOPE_VERSION_COLUMN => {
912 if let JsonValue::String(scope_version) = cell
913 && !scope_version.is_empty()
914 {
915 span.tags.push(KeyValue {
916 key: KEY_OTEL_SCOPE_VERSION.to_string(),
917 value_type: ValueType::String,
918 value: Value::String(scope_version),
919 });
920 }
921 }
922 SPAN_KIND_COLUMN => {
923 if let JsonValue::String(span_kind) = cell
924 && !span_kind.is_empty()
925 {
926 span.tags.push(KeyValue {
927 key: KEY_SPAN_KIND.to_string(),
928 value_type: ValueType::String,
929 value: Value::String(normalize_span_kind(&span_kind)),
930 });
931 }
932 }
933 SPAN_STATUS_CODE => {
934 if let JsonValue::String(span_status) = cell
935 && span_status != SPAN_STATUS_UNSET
936 && !span_status.is_empty()
937 {
938 span.tags.push(KeyValue {
939 key: KEY_OTEL_STATUS_CODE.to_string(),
940 value_type: ValueType::String,
941 value: Value::String(normalize_status_code(&span_status)),
942 });
943 if span_status == SPAN_STATUS_ERROR {
945 span.tags.push(KeyValue {
946 key: KEY_OTEL_STATUS_ERROR_KEY.to_string(),
947 value_type: ValueType::Boolean,
948 value: Value::Boolean(true),
949 });
950 }
951 }
952 }
953
954 SPAN_STATUS_MESSAGE_COLUMN => {
955 if let JsonValue::String(span_status_message) = cell
956 && !span_status_message.is_empty()
957 {
958 span.tags.push(KeyValue {
959 key: KEY_OTEL_STATUS_MESSAGE.to_string(),
960 value_type: ValueType::String,
961 value: Value::String(span_status_message),
962 });
963 }
964 }
965
966 TRACE_STATE_COLUMN => {
967 if let JsonValue::String(trace_state) = cell
968 && !trace_state.is_empty()
969 {
970 span.tags.push(KeyValue {
971 key: KEY_OTEL_TRACE_STATE.to_string(),
972 value_type: ValueType::String,
973 value: Value::String(trace_state),
974 });
975 }
976 }
977
978 _ => {
979 if is_span_attributes_flatten {
981 const SPAN_ATTR_PREFIX: &str = "span_attributes.";
982 const RESOURCE_ATTR_PREFIX: &str = "resource_attributes.";
983 if column_name.starts_with(SPAN_ATTR_PREFIX) {
985 if let Some(keyvalue) = to_keyvalue(
986 column_name
987 .strip_prefix(SPAN_ATTR_PREFIX)
988 .unwrap_or_default()
989 .to_string(),
990 cell,
991 ) {
992 span.tags.push(keyvalue);
993 }
994 } else if column_name.starts_with(RESOURCE_ATTR_PREFIX)
995 && let Some(keyvalue) = to_keyvalue(
996 column_name
997 .strip_prefix(RESOURCE_ATTR_PREFIX)
998 .unwrap_or_default()
999 .to_string(),
1000 cell,
1001 )
1002 {
1003 resource_tags.push(keyvalue);
1004 }
1005 }
1006 }
1007 }
1008 }
1009
1010 if let Some(parent_span_id) = parent_span_id {
1011 span.references.insert(
1012 0,
1013 Reference {
1014 trace_id: span.trace_id.clone(),
1015 span_id: parent_span_id,
1016 ref_type: REF_TYPE_CHILD_OF.to_string(),
1017 },
1018 );
1019 }
1020
1021 if let Some(service_name) = service_name {
1022 if !service_to_resource_attributes.contains_key(&service_name) {
1023 service_to_resource_attributes.insert(service_name.clone(), resource_tags);
1024 }
1025
1026 if let Some(process) = trace_id_to_processes.get_mut(&span.trace_id) {
1027 if let Some(process_id) = process.get(&service_name) {
1028 span.process_id = process_id.clone();
1029 } else {
1030 let process_id = format!("p{}", process.len() + 1);
1032 process.insert(service_name, process_id.clone());
1033 span.process_id = process_id;
1034 }
1035 }
1036 }
1037
1038 span.tags.sort_by(|a, b| a.key.cmp(&b.key));
1040
1041 if let Some(spans) = trace_id_to_spans.get_mut(&span.trace_id) {
1042 spans.push(span);
1043 } else {
1044 trace_id_to_spans.insert(span.trace_id.clone(), vec![span]);
1045 }
1046 }
1047
1048 let mut traces = Vec::new();
1049 for (trace_id, spans) in trace_id_to_spans {
1050 let mut trace = Trace {
1051 trace_id,
1052 spans,
1053 ..Default::default()
1054 };
1055
1056 if let Some(processes) = trace_id_to_processes.remove(&trace.trace_id) {
1057 let mut process_id_to_process = HashMap::new();
1058 for (service_name, process_id) in processes.into_iter() {
1059 let tags = service_to_resource_attributes
1060 .remove(&service_name)
1061 .unwrap_or_default();
1062 process_id_to_process.insert(process_id, Process { service_name, tags });
1063 }
1064 trace.processes = process_id_to_process;
1065 }
1066 traces.push(trace);
1067 }
1068
1069 Ok(traces)
1070}
1071
1072fn to_keyvalue(key: String, value: JsonValue) -> Option<KeyValue> {
1073 match value {
1074 JsonValue::String(value) => Some(KeyValue {
1075 key,
1076 value_type: ValueType::String,
1077 value: Value::String(value.clone()),
1078 }),
1079 JsonValue::Number(value) => {
1080 if value.is_i64() {
1081 Some(KeyValue {
1082 key,
1083 value_type: ValueType::Int64,
1084 value: Value::Int64(value.as_i64()?),
1085 })
1086 } else {
1087 Some(KeyValue {
1088 key,
1089 value_type: ValueType::Float64,
1090 value: Value::Float64(value.as_f64()?),
1091 })
1092 }
1093 }
1094 JsonValue::Bool(value) => Some(KeyValue {
1095 key,
1096 value_type: ValueType::Boolean,
1097 value: Value::Boolean(value),
1098 }),
1099 JsonValue::Array(value) => Some(KeyValue {
1100 key,
1101 value_type: ValueType::String,
1102 value: Value::String(serde_json::to_string(&value).unwrap()),
1103 }),
1104 JsonValue::Object(value) => Some(KeyValue {
1105 key,
1106 value_type: ValueType::String,
1107 value: Value::String(serde_json::to_string(&value).unwrap()),
1108 }),
1109 JsonValue::Null => None,
1110 }
1111}
1112
1113fn object_to_tags(object: serde_json::map::Map<String, JsonValue>) -> Vec<KeyValue> {
1114 object
1115 .into_iter()
1116 .filter_map(|(key, value)| to_keyvalue(key, value))
1117 .collect()
1118}
1119
1120fn services_from_records(records: HttpRecordsOutput) -> Result<Vec<String>> {
1121 let expected_schema = vec![(SERVICE_NAME_COLUMN, "String")];
1122 check_schema(&records, &expected_schema)?;
1123
1124 let mut services = Vec::with_capacity(records.total_rows);
1125 for row in records.rows.into_iter() {
1126 for value in row.into_iter() {
1127 if let JsonValue::String(service_name) = value {
1128 services.push(service_name);
1129 }
1130 }
1131 }
1132 Ok(services)
1133}
1134
1135fn operations_from_records(
1137 records: HttpRecordsOutput,
1138 contain_span_kind: bool,
1139) -> Result<Vec<Operation>> {
1140 let expected_schema = vec![(SPAN_NAME_COLUMN, "String"), (SPAN_KIND_COLUMN, "String")];
1141 check_schema(&records, &expected_schema)?;
1142
1143 let mut operations = Vec::with_capacity(records.total_rows);
1144 for row in records.rows.into_iter() {
1145 let mut row_iter = row.into_iter();
1146 if let Some(JsonValue::String(operation)) = row_iter.next() {
1147 let mut operation = Operation {
1148 name: operation,
1149 span_kind: None,
1150 };
1151 if contain_span_kind {
1152 if let Some(JsonValue::String(span_kind)) = row_iter.next() {
1153 operation.span_kind = Some(normalize_span_kind(&span_kind));
1154 }
1155 } else {
1156 row_iter.next();
1158 }
1159 operations.push(operation);
1160 }
1161 }
1162
1163 Ok(operations)
1164}
1165
1166fn check_schema(records: &HttpRecordsOutput, expected_schema: &[(&str, &str)]) -> Result<()> {
1168 for (i, column) in records.schema.column_schemas.iter().enumerate() {
1169 if column.name != expected_schema[i].0 || column.data_type != expected_schema[i].1 {
1170 InvalidJaegerQuerySnafu {
1171 reason: "query result schema is not correct".to_string(),
1172 }
1173 .fail()?
1174 }
1175 }
1176 Ok(())
1177}
1178
1179fn normalize_span_kind(span_kind: &str) -> String {
1182 if let Some(stripped) = span_kind.strip_prefix(SPAN_KIND_PREFIX) {
1184 stripped.to_lowercase()
1185 } else {
1186 span_kind.to_lowercase()
1188 }
1189}
1190
1191fn normalize_status_code(status_code: &str) -> String {
1194 if let Some(stripped) = status_code.strip_prefix(SPAN_STATUS_PREFIX) {
1196 stripped.to_string()
1197 } else {
1198 status_code.to_string()
1200 }
1201}
1202
1203fn convert_string_to_number(input: &serde_json::Value) -> Option<serde_json::Value> {
1204 if let Some(data) = input.as_str() {
1205 if let Ok(number) = data.parse::<i64>() {
1206 return Some(serde_json::Value::Number(serde_json::Number::from(number)));
1207 }
1208 if let Ok(number) = data.parse::<f64>()
1209 && let Some(number) = serde_json::Number::from_f64(number)
1210 {
1211 return Some(serde_json::Value::Number(number));
1212 }
1213 }
1214
1215 None
1216}
1217
1218fn convert_string_to_boolean(input: &serde_json::Value) -> Option<serde_json::Value> {
1219 if let Some(data) = input.as_str() {
1220 if data == "true" {
1221 return Some(serde_json::Value::Bool(true));
1222 }
1223 if data == "false" {
1224 return Some(serde_json::Value::Bool(false));
1225 }
1226 }
1227
1228 None
1229}
1230
1231#[cfg(test)]
1232mod tests {
1233 use serde_json::{Number, Value as JsonValue, json};
1234
1235 use super::*;
1236 use crate::http::{ColumnSchema, HttpRecordsOutput, OutputSchema};
1237
1238 #[test]
1239 fn test_numeric_tag_types() {
1240 for (input, value_type, value) in [
1241 (json!(42), ValueType::Int64, Value::Int64(42)),
1242 (json!(1.5), ValueType::Float64, Value::Float64(1.5)),
1243 ] {
1244 assert_eq!(
1245 to_keyvalue("tag".to_string(), input),
1246 Some(KeyValue {
1247 key: "tag".to_string(),
1248 value_type,
1249 value,
1250 })
1251 );
1252 }
1253 }
1254
1255 #[test]
1256 fn test_services_from_records() {
1257 let tests = vec![(
1259 HttpRecordsOutput {
1260 schema: OutputSchema {
1261 column_schemas: vec![ColumnSchema {
1262 name: "service_name".to_string(),
1263 data_type: "String".to_string(),
1264 }],
1265 },
1266 rows: vec![
1267 vec![JsonValue::String("test-service-0".to_string())],
1268 vec![JsonValue::String("test-service-1".to_string())],
1269 ],
1270 total_rows: 2,
1271 metrics: HashMap::new(),
1272 },
1273 vec!["test-service-0".to_string(), "test-service-1".to_string()],
1274 )];
1275
1276 for (records, expected) in tests {
1277 let services = services_from_records(records).unwrap();
1278 assert_eq!(services, expected);
1279 }
1280 }
1281
1282 #[test]
1283 fn test_operations_from_records() {
1284 let tests = vec![
1286 (
1287 HttpRecordsOutput {
1288 schema: OutputSchema {
1289 column_schemas: vec![
1290 ColumnSchema {
1291 name: "span_name".to_string(),
1292 data_type: "String".to_string(),
1293 },
1294 ColumnSchema {
1295 name: "span_kind".to_string(),
1296 data_type: "String".to_string(),
1297 },
1298 ],
1299 },
1300 rows: vec![
1301 vec![
1302 JsonValue::String("access-mysql".to_string()),
1303 JsonValue::String("SPAN_KIND_SERVER".to_string()),
1304 ],
1305 vec![
1306 JsonValue::String("access-redis".to_string()),
1307 JsonValue::String("SPAN_KIND_CLIENT".to_string()),
1308 ],
1309 ],
1310 total_rows: 2,
1311 metrics: HashMap::new(),
1312 },
1313 false,
1314 vec![
1315 Operation {
1316 name: "access-mysql".to_string(),
1317 span_kind: None,
1318 },
1319 Operation {
1320 name: "access-redis".to_string(),
1321 span_kind: None,
1322 },
1323 ],
1324 ),
1325 (
1326 HttpRecordsOutput {
1327 schema: OutputSchema {
1328 column_schemas: vec![
1329 ColumnSchema {
1330 name: "span_name".to_string(),
1331 data_type: "String".to_string(),
1332 },
1333 ColumnSchema {
1334 name: "span_kind".to_string(),
1335 data_type: "String".to_string(),
1336 },
1337 ],
1338 },
1339 rows: vec![
1340 vec![
1341 JsonValue::String("access-mysql".to_string()),
1342 JsonValue::String("SPAN_KIND_SERVER".to_string()),
1343 ],
1344 vec![
1345 JsonValue::String("access-redis".to_string()),
1346 JsonValue::String("SPAN_KIND_CLIENT".to_string()),
1347 ],
1348 ],
1349 total_rows: 2,
1350 metrics: HashMap::new(),
1351 },
1352 true,
1353 vec![
1354 Operation {
1355 name: "access-mysql".to_string(),
1356 span_kind: Some("server".to_string()),
1357 },
1358 Operation {
1359 name: "access-redis".to_string(),
1360 span_kind: Some("client".to_string()),
1361 },
1362 ],
1363 ),
1364 ];
1365
1366 for (records, contain_span_kind, expected) in tests {
1367 let operations = operations_from_records(records, contain_span_kind).unwrap();
1368 assert_eq!(operations, expected);
1369 }
1370 }
1371
1372 #[test]
1373 fn test_traces_from_records() {
1374 let tests = vec![(
1376 HttpRecordsOutput {
1377 schema: OutputSchema {
1378 column_schemas: vec![
1379 ColumnSchema {
1380 name: "trace_id".to_string(),
1381 data_type: "String".to_string(),
1382 },
1383 ColumnSchema {
1384 name: "timestamp".to_string(),
1385 data_type: "TimestampNanosecond".to_string(),
1386 },
1387 ColumnSchema {
1388 name: "duration_nano".to_string(),
1389 data_type: "UInt64".to_string(),
1390 },
1391 ColumnSchema {
1392 name: "service_name".to_string(),
1393 data_type: "String".to_string(),
1394 },
1395 ColumnSchema {
1396 name: "span_name".to_string(),
1397 data_type: "String".to_string(),
1398 },
1399 ColumnSchema {
1400 name: "span_id".to_string(),
1401 data_type: "String".to_string(),
1402 },
1403 ColumnSchema {
1404 name: "span_attributes".to_string(),
1405 data_type: "Json".to_string(),
1406 },
1407 ],
1408 },
1409 rows: vec![
1410 vec![
1411 JsonValue::String("5611dce1bc9ebed65352d99a027b08ea".to_string()),
1412 JsonValue::Number(Number::from_u128(1738726754492422000).unwrap()),
1413 JsonValue::Number(Number::from_u128(100000000).unwrap()),
1414 JsonValue::String("test-service-0".to_string()),
1415 JsonValue::String("access-mysql".to_string()),
1416 JsonValue::String("008421dbbd33a3e9".to_string()),
1417 JsonValue::Object(
1418 json!({
1419 "operation.type": "access-mysql",
1420 })
1421 .as_object()
1422 .unwrap()
1423 .clone(),
1424 ),
1425 ],
1426 vec![
1427 JsonValue::String("5611dce1bc9ebed65352d99a027b08ea".to_string()),
1428 JsonValue::Number(Number::from_u128(1738726754642422000).unwrap()),
1429 JsonValue::Number(Number::from_u128(100000000).unwrap()),
1430 JsonValue::String("test-service-0".to_string()),
1431 JsonValue::String("access-redis".to_string()),
1432 JsonValue::String("ffa03416a7b9ea48".to_string()),
1433 JsonValue::Object(
1434 json!({
1435 "operation.type": "access-redis",
1436 })
1437 .as_object()
1438 .unwrap()
1439 .clone(),
1440 ),
1441 ],
1442 ],
1443 total_rows: 2,
1444 metrics: HashMap::new(),
1445 },
1446 vec![Trace {
1447 trace_id: "5611dce1bc9ebed65352d99a027b08ea".to_string(),
1448 spans: vec![
1449 Span {
1450 trace_id: "5611dce1bc9ebed65352d99a027b08ea".to_string(),
1451 span_id: "008421dbbd33a3e9".to_string(),
1452 operation_name: "access-mysql".to_string(),
1453 start_time: 1738726754492422,
1454 duration: 100000,
1455 tags: vec![KeyValue {
1456 key: "operation.type".to_string(),
1457 value_type: ValueType::String,
1458 value: Value::String("access-mysql".to_string()),
1459 }],
1460 process_id: "p1".to_string(),
1461 ..Default::default()
1462 },
1463 Span {
1464 trace_id: "5611dce1bc9ebed65352d99a027b08ea".to_string(),
1465 span_id: "ffa03416a7b9ea48".to_string(),
1466 operation_name: "access-redis".to_string(),
1467 start_time: 1738726754642422,
1468 duration: 100000,
1469 tags: vec![KeyValue {
1470 key: "operation.type".to_string(),
1471 value_type: ValueType::String,
1472 value: Value::String("access-redis".to_string()),
1473 }],
1474 process_id: "p1".to_string(),
1475 ..Default::default()
1476 },
1477 ],
1478 processes: HashMap::from([(
1479 "p1".to_string(),
1480 Process {
1481 service_name: "test-service-0".to_string(),
1482 tags: vec![],
1483 },
1484 )]),
1485 ..Default::default()
1486 }],
1487 )];
1488
1489 for (records, expected) in tests {
1490 let traces = traces_from_records(records).unwrap();
1491 assert_eq!(traces, expected);
1492 }
1493 }
1494
1495 #[test]
1496 fn test_traces_from_v1_records() {
1497 let tests = vec![(
1499 HttpRecordsOutput {
1500 schema: OutputSchema {
1501 column_schemas: vec![
1502 ColumnSchema {
1503 name: "trace_id".to_string(),
1504 data_type: "String".to_string(),
1505 },
1506 ColumnSchema {
1507 name: "timestamp".to_string(),
1508 data_type: "TimestampNanosecond".to_string(),
1509 },
1510 ColumnSchema {
1511 name: "duration_nano".to_string(),
1512 data_type: "UInt64".to_string(),
1513 },
1514 ColumnSchema {
1515 name: "service_name".to_string(),
1516 data_type: "String".to_string(),
1517 },
1518 ColumnSchema {
1519 name: "span_name".to_string(),
1520 data_type: "String".to_string(),
1521 },
1522 ColumnSchema {
1523 name: "span_id".to_string(),
1524 data_type: "String".to_string(),
1525 },
1526 ColumnSchema {
1527 name: "span_attributes.http.request.method".to_string(),
1528 data_type: "String".to_string(),
1529 },
1530 ColumnSchema {
1531 name: "span_attributes.http.request.url".to_string(),
1532 data_type: "String".to_string(),
1533 },
1534 ColumnSchema {
1535 name: "span_attributes.http.status_code".to_string(),
1536 data_type: "UInt64".to_string(),
1537 },
1538 ],
1539 },
1540 rows: vec![
1541 vec![
1542 JsonValue::String("5611dce1bc9ebed65352d99a027b08ea".to_string()),
1543 JsonValue::Number(Number::from_u128(1738726754492422000).unwrap()),
1544 JsonValue::Number(Number::from_u128(100000000).unwrap()),
1545 JsonValue::String("test-service-0".to_string()),
1546 JsonValue::String("access-mysql".to_string()),
1547 JsonValue::String("008421dbbd33a3e9".to_string()),
1548 JsonValue::String("GET".to_string()),
1549 JsonValue::String("/data".to_string()),
1550 JsonValue::Number(Number::from_u128(200).unwrap()),
1551 ],
1552 vec![
1553 JsonValue::String("5611dce1bc9ebed65352d99a027b08ea".to_string()),
1554 JsonValue::Number(Number::from_u128(1738726754642422000).unwrap()),
1555 JsonValue::Number(Number::from_u128(100000000).unwrap()),
1556 JsonValue::String("test-service-0".to_string()),
1557 JsonValue::String("access-redis".to_string()),
1558 JsonValue::String("ffa03416a7b9ea48".to_string()),
1559 JsonValue::String("POST".to_string()),
1560 JsonValue::String("/create".to_string()),
1561 JsonValue::Number(Number::from_u128(400).unwrap()),
1562 ],
1563 ],
1564 total_rows: 2,
1565 metrics: HashMap::new(),
1566 },
1567 vec![Trace {
1568 trace_id: "5611dce1bc9ebed65352d99a027b08ea".to_string(),
1569 spans: vec![
1570 Span {
1571 trace_id: "5611dce1bc9ebed65352d99a027b08ea".to_string(),
1572 span_id: "008421dbbd33a3e9".to_string(),
1573 operation_name: "access-mysql".to_string(),
1574 start_time: 1738726754492422,
1575 duration: 100000,
1576 tags: vec![
1577 KeyValue {
1578 key: "http.request.method".to_string(),
1579 value_type: ValueType::String,
1580 value: Value::String("GET".to_string()),
1581 },
1582 KeyValue {
1583 key: "http.request.url".to_string(),
1584 value_type: ValueType::String,
1585 value: Value::String("/data".to_string()),
1586 },
1587 KeyValue {
1588 key: "http.status_code".to_string(),
1589 value_type: ValueType::Int64,
1590 value: Value::Int64(200),
1591 },
1592 ],
1593 process_id: "p1".to_string(),
1594 ..Default::default()
1595 },
1596 Span {
1597 trace_id: "5611dce1bc9ebed65352d99a027b08ea".to_string(),
1598 span_id: "ffa03416a7b9ea48".to_string(),
1599 operation_name: "access-redis".to_string(),
1600 start_time: 1738726754642422,
1601 duration: 100000,
1602 tags: vec![
1603 KeyValue {
1604 key: "http.request.method".to_string(),
1605 value_type: ValueType::String,
1606 value: Value::String("POST".to_string()),
1607 },
1608 KeyValue {
1609 key: "http.request.url".to_string(),
1610 value_type: ValueType::String,
1611 value: Value::String("/create".to_string()),
1612 },
1613 KeyValue {
1614 key: "http.status_code".to_string(),
1615 value_type: ValueType::Int64,
1616 value: Value::Int64(400),
1617 },
1618 ],
1619 process_id: "p1".to_string(),
1620 ..Default::default()
1621 },
1622 ],
1623 processes: HashMap::from([(
1624 "p1".to_string(),
1625 Process {
1626 service_name: "test-service-0".to_string(),
1627 tags: vec![],
1628 },
1629 )]),
1630 ..Default::default()
1631 }],
1632 )];
1633
1634 for (records, expected) in tests {
1635 let traces = traces_from_records(records).unwrap();
1636 assert_eq!(traces, expected);
1637 }
1638 }
1639
1640 #[test]
1641 fn test_from_jaeger_query_params() {
1642 let tests = vec![
1644 (
1645 JaegerQueryParams {
1646 service_name: Some("test-service-0".to_string()),
1647 ..Default::default()
1648 },
1649 QueryTraceParams {
1650 service_name: "test-service-0".to_string(),
1651 ..Default::default()
1652 },
1653 ),
1654 (
1655 JaegerQueryParams {
1656 service_name: Some("test-service-0".to_string()),
1657 operation_name: Some("access-mysql".to_string()),
1658 start: Some(1738726754492422),
1659 end: Some(1738726754642422),
1660 max_duration: Some("100ms".to_string()),
1661 min_duration: Some("50ms".to_string()),
1662 limit: Some(10),
1663 tags: Some("{\"http.status_code\":\"200\",\"latency\":\"11.234\",\"error\":\"false\",\"http.method\":\"GET\",\"http.path\":\"/api/v1/users\"}".to_string()),
1664 ..Default::default()
1665 },
1666 QueryTraceParams {
1667 service_name: "test-service-0".to_string(),
1668 operation_name: Some("access-mysql".to_string()),
1669 start_time: Some(1738726754492422000),
1670 end_time: Some(1738726754642422000),
1671 min_duration: Some(50000000),
1672 max_duration: Some(100000000),
1673 limit: Some(10),
1674 tags: Some(HashMap::from([
1675 ("http.status_code".to_string(), JsonValue::Number(Number::from(200))),
1676 ("latency".to_string(), JsonValue::Number(Number::from_f64(11.234).unwrap())),
1677 ("error".to_string(), JsonValue::Bool(false)),
1678 ("http.method".to_string(), JsonValue::String("GET".to_string())),
1679 ("http.path".to_string(), JsonValue::String("/api/v1/users".to_string())),
1680 ])),
1681 user_agent: TraceUserAgent::Jaeger,
1682 },
1683 ),
1684 ];
1685
1686 for (query_params, expected) in tests {
1687 let query_params = QueryTraceParams::from_jaeger_query_params(query_params).unwrap();
1688 assert_eq!(query_params, expected);
1689 }
1690 }
1691
1692 #[test]
1693 fn test_check_schema() {
1694 let tests = vec![(
1696 HttpRecordsOutput {
1697 schema: OutputSchema {
1698 column_schemas: vec![
1699 ColumnSchema {
1700 name: "trace_id".to_string(),
1701 data_type: "String".to_string(),
1702 },
1703 ColumnSchema {
1704 name: "timestamp".to_string(),
1705 data_type: "TimestampNanosecond".to_string(),
1706 },
1707 ColumnSchema {
1708 name: "duration_nano".to_string(),
1709 data_type: "UInt64".to_string(),
1710 },
1711 ColumnSchema {
1712 name: "service_name".to_string(),
1713 data_type: "String".to_string(),
1714 },
1715 ColumnSchema {
1716 name: "span_name".to_string(),
1717 data_type: "String".to_string(),
1718 },
1719 ColumnSchema {
1720 name: "span_id".to_string(),
1721 data_type: "String".to_string(),
1722 },
1723 ColumnSchema {
1724 name: "span_attributes".to_string(),
1725 data_type: "Json".to_string(),
1726 },
1727 ],
1728 },
1729 rows: vec![],
1730 total_rows: 0,
1731 metrics: HashMap::new(),
1732 },
1733 vec![
1734 (TRACE_ID_COLUMN, "String"),
1735 (TIMESTAMP_COLUMN, "TimestampNanosecond"),
1736 (DURATION_NANO_COLUMN, "UInt64"),
1737 (SERVICE_NAME_COLUMN, "String"),
1738 (SPAN_NAME_COLUMN, "String"),
1739 (SPAN_ID_COLUMN, "String"),
1740 (SPAN_ATTRIBUTES_COLUMN, "Json"),
1741 ],
1742 true,
1743 )];
1744
1745 for (records, expected_schema, is_ok) in tests {
1746 let result = check_schema(&records, &expected_schema);
1747 assert_eq!(result.is_ok(), is_ok);
1748 }
1749 }
1750
1751 #[test]
1752 fn test_normalize_span_kind() {
1753 let tests = vec![
1754 ("SPAN_KIND_SERVER".to_string(), "server".to_string()),
1755 ("SPAN_KIND_CLIENT".to_string(), "client".to_string()),
1756 ];
1757
1758 for (input, expected) in tests {
1759 let result = normalize_span_kind(&input);
1760 assert_eq!(result, expected);
1761 }
1762 }
1763
1764 #[test]
1765 fn test_convert_string_to_number() {
1766 let tests = vec![
1767 (
1768 JsonValue::String("123".to_string()),
1769 Some(JsonValue::Number(Number::from(123))),
1770 ),
1771 (
1772 JsonValue::String("123.456".to_string()),
1773 Some(JsonValue::Number(Number::from_f64(123.456).unwrap())),
1774 ),
1775 ];
1776
1777 for (input, expected) in tests {
1778 let result = convert_string_to_number(&input);
1779 assert_eq!(result, expected);
1780 }
1781 }
1782
1783 #[test]
1784 fn test_convert_string_to_boolean() {
1785 let tests = vec![
1786 (
1787 JsonValue::String("true".to_string()),
1788 Some(JsonValue::Bool(true)),
1789 ),
1790 (
1791 JsonValue::String("false".to_string()),
1792 Some(JsonValue::Bool(false)),
1793 ),
1794 ];
1795
1796 for (input, expected) in tests {
1797 let result = convert_string_to_boolean(&input);
1798 assert_eq!(result, expected);
1799 }
1800 }
1801}