1#![allow(dead_code)]
2use std::collections::HashMap;
31use std::time::{Duration, SystemTime, UNIX_EPOCH};
32
33use serde::{Deserialize, Serialize};
34use thiserror::Error;
35use tracing::{debug, info, warn};
36
37#[derive(Debug, Error)]
41pub enum TraceExportError {
42 #[error("backend endpoint not configured: {0}")]
44 NotConfigured(String),
45 #[error("HTTP export failed: {0}")]
47 Http(String),
48 #[error("serialisation error: {0}")]
50 Serialisation(String),
51 #[error("internal lock poisoned")]
53 LockPoisoned,
54}
55
56#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
60pub enum PipelineStage {
61 Rag,
62 Assemble,
63 Inference,
64 Post,
65 Stream,
66}
67
68impl PipelineStage {
69 #[must_use]
71 pub fn name(self) -> &'static str {
72 match self {
73 PipelineStage::Rag => "rag",
74 PipelineStage::Assemble => "assemble",
75 PipelineStage::Inference => "inference",
76 PipelineStage::Post => "post",
77 PipelineStage::Stream => "stream",
78 }
79 }
80
81 #[must_use]
83 pub fn index(self) -> usize {
84 match self {
85 PipelineStage::Rag => 0,
86 PipelineStage::Assemble => 1,
87 PipelineStage::Inference => 2,
88 PipelineStage::Post => 3,
89 PipelineStage::Stream => 4,
90 }
91 }
92}
93
94#[derive(Debug, Clone, Serialize, Deserialize)]
96pub struct PipelineSpan {
97 pub request_id: String,
99 pub session_id: String,
101 pub stage: PipelineStage,
103 pub start_unix_nanos: u128,
105 pub duration: Duration,
107 pub success: bool,
109 pub error: Option<String>,
111 pub tags: HashMap<String, String>,
113}
114
115impl PipelineSpan {
116 #[must_use]
118 pub fn start(
119 request_id: impl Into<String>,
120 session_id: impl Into<String>,
121 stage: PipelineStage,
122 ) -> Self {
123 let start_unix_nanos = SystemTime::now()
124 .duration_since(UNIX_EPOCH)
125 .unwrap_or_default()
126 .as_nanos();
127 Self {
128 request_id: request_id.into(),
129 session_id: session_id.into(),
130 stage,
131 start_unix_nanos,
132 duration: Duration::ZERO,
133 success: false,
134 error: None,
135 tags: HashMap::new(),
136 }
137 }
138
139 pub fn finish(&mut self, success: bool, error: Option<String>) {
141 let now_nanos = SystemTime::now()
142 .duration_since(UNIX_EPOCH)
143 .unwrap_or_default()
144 .as_nanos();
145 let elapsed_nanos = now_nanos.saturating_sub(self.start_unix_nanos);
146 self.duration = Duration::from_nanos(elapsed_nanos as u64);
147 self.success = success;
148 self.error = error;
149 }
150
151 pub fn tag(&mut self, key: impl Into<String>, value: impl Into<String>) {
153 self.tags.insert(key.into(), value.into());
154 }
155}
156
157#[derive(Debug, Clone, Serialize, Deserialize)]
161pub struct RequestTrace {
162 pub request_id: String,
164 pub session_id: String,
166 pub entry_unix_nanos: u128,
168 pub total_duration: Duration,
170 pub success: bool,
172 pub spans: Vec<PipelineSpan>,
174 pub service_name: String,
176}
177
178impl RequestTrace {
179 #[must_use]
181 pub fn new(
182 request_id: impl Into<String>,
183 session_id: impl Into<String>,
184 service_name: impl Into<String>,
185 ) -> Self {
186 let entry_unix_nanos = SystemTime::now()
187 .duration_since(UNIX_EPOCH)
188 .unwrap_or_default()
189 .as_nanos();
190 Self {
191 request_id: request_id.into(),
192 session_id: session_id.into(),
193 entry_unix_nanos,
194 total_duration: Duration::ZERO,
195 success: false,
196 spans: Vec::new(),
197 service_name: service_name.into(),
198 }
199 }
200
201 pub fn add_span(&mut self, span: PipelineSpan) {
203 self.spans.push(span);
204 }
205
206 pub fn finish(&mut self, success: bool) {
208 let now_nanos = SystemTime::now()
209 .duration_since(UNIX_EPOCH)
210 .unwrap_or_default()
211 .as_nanos();
212 let elapsed_nanos = now_nanos.saturating_sub(self.entry_unix_nanos);
213 self.total_duration = Duration::from_nanos(elapsed_nanos as u64);
214 self.success = success;
215 }
216}
217
218#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
222pub enum TraceBackend {
223 Jaeger,
225 #[cfg(feature = "zipkin")]
227 Zipkin,
228 #[cfg(feature = "datadog")]
230 DataDog,
231 NoOp,
233}
234
235#[derive(Debug, Clone, Serialize, Deserialize)]
239pub struct TraceExporterConfig {
240 pub backend: TraceBackend,
242 pub service_name: String,
244 pub endpoint_override: Option<String>,
246 pub export_timeout: Duration,
248}
249
250impl Default for TraceExporterConfig {
251 fn default() -> Self {
252 Self {
253 backend: TraceBackend::Jaeger,
254 service_name: "tokio-prompt-orchestrator".to_string(),
255 endpoint_override: None,
256 export_timeout: Duration::from_secs(5),
257 }
258 }
259}
260
261impl TraceExporterConfig {
262 #[must_use]
264 pub fn from_env() -> Self {
265 let backend = {
266 #[cfg(feature = "datadog")]
267 if std::env::var("DD_AGENT_HOST").is_ok() {
268 TraceBackend::DataDog
269 } else {
270 #[cfg(feature = "zipkin")]
271 if std::env::var("ZIPKIN_ENDPOINT").is_ok() {
272 TraceBackend::Zipkin
273 } else {
274 TraceBackend::Jaeger
275 }
276 #[cfg(not(feature = "zipkin"))]
277 TraceBackend::Jaeger
278 }
279 #[cfg(not(feature = "datadog"))]
280 {
281 #[cfg(feature = "zipkin")]
282 if std::env::var("ZIPKIN_ENDPOINT").is_ok() {
283 TraceBackend::Zipkin
284 } else {
285 TraceBackend::Jaeger
286 }
287 #[cfg(not(feature = "zipkin"))]
288 TraceBackend::Jaeger
289 }
290 };
291 Self {
292 backend,
293 ..Default::default()
294 }
295 }
296
297 #[must_use]
299 pub fn resolved_endpoint(&self) -> Option<String> {
300 if let Some(ref ep) = self.endpoint_override {
301 return Some(ep.clone());
302 }
303 match &self.backend {
304 TraceBackend::Jaeger => std::env::var("OTEL_EXPORTER_OTLP_ENDPOINT")
305 .or_else(|_| std::env::var("JAEGER_ENDPOINT"))
306 .ok(),
307 #[cfg(feature = "zipkin")]
308 TraceBackend::Zipkin => std::env::var("ZIPKIN_ENDPOINT").ok(),
309 #[cfg(feature = "datadog")]
310 TraceBackend::DataDog => {
311 let host = std::env::var("DD_AGENT_HOST").unwrap_or_else(|_| "localhost".into());
312 let port = std::env::var("DD_AGENT_PORT").unwrap_or_else(|_| "8126".into());
313 Some(format!("http://{host}:{port}"))
314 }
315 TraceBackend::NoOp => None,
316 }
317 }
318}
319
320pub struct TraceExporter {
331 config: TraceExporterConfig,
332 http_client: reqwest::Client,
333}
334
335impl TraceExporter {
336 pub fn new(config: TraceExporterConfig) -> Result<Self, TraceExportError> {
341 let http_client = reqwest::Client::builder()
342 .timeout(config.export_timeout)
343 .build()
344 .map_err(|e| TraceExportError::Http(e.to_string()))?;
345 Ok(Self { config, http_client })
346 }
347
348 pub async fn export(&self, trace: &RequestTrace) -> Result<(), TraceExportError> {
353 match &self.config.backend {
354 TraceBackend::NoOp => {
355 debug!(request_id = %trace.request_id, "TraceExporter: no-op export");
356 return Ok(());
357 }
358 TraceBackend::Jaeger => self.export_otlp(trace).await?,
359 #[cfg(feature = "zipkin")]
360 TraceBackend::Zipkin => self.export_zipkin(trace).await?,
361 #[cfg(feature = "datadog")]
362 TraceBackend::DataDog => self.export_datadog(trace).await?,
363 }
364 info!(
365 request_id = %trace.request_id,
366 backend = ?self.config.backend,
367 spans = trace.spans.len(),
368 duration_ms = trace.total_duration.as_millis(),
369 "trace exported"
370 );
371 Ok(())
372 }
373
374 async fn export_otlp(&self, trace: &RequestTrace) -> Result<(), TraceExportError> {
377 let endpoint = self
378 .config
379 .resolved_endpoint()
380 .ok_or_else(|| TraceExportError::NotConfigured("OTEL_EXPORTER_OTLP_ENDPOINT or JAEGER_ENDPOINT".to_string()))?;
381
382 let payload = self.build_otlp_payload(trace)?;
383 let url = format!("{endpoint}/v1/traces");
384
385 let resp = self
386 .http_client
387 .post(&url)
388 .header("Content-Type", "application/json")
389 .body(payload)
390 .send()
391 .await
392 .map_err(|e| TraceExportError::Http(e.to_string()))?;
393
394 if !resp.status().is_success() {
395 warn!(
396 status = %resp.status(),
397 url = %url,
398 "OTLP trace export received non-2xx response"
399 );
400 }
401 Ok(())
402 }
403
404 fn build_otlp_payload(&self, trace: &RequestTrace) -> Result<String, TraceExportError> {
405 let trace_id = hex_from_str(&trace.request_id, 32);
407 let root_span_id = hex_from_str(&trace.request_id, 16);
408
409 let mut span_jsons: Vec<String> = Vec::with_capacity(trace.spans.len() + 1);
410
411 let root_end_nanos = trace.entry_unix_nanos
413 + trace.total_duration.as_nanos();
414 span_jsons.push(format!(
415 r#"{{"traceId":"{trace_id}","spanId":"{root_span_id}","name":"request","kind":1,"startTimeUnixNano":"{start}","endTimeUnixNano":"{end}","status":{{"code":{code}}},"attributes":[{{"key":"request_id","value":{{"stringValue":"{rid}"}}}},{{"key":"session_id","value":{{"stringValue":"{sid}"}}}}]}}"#,
416 trace_id = trace_id,
417 root_span_id = root_span_id,
418 start = trace.entry_unix_nanos,
419 end = root_end_nanos,
420 code = if trace.success { 1 } else { 2 },
421 rid = trace.request_id,
422 sid = trace.session_id,
423 ));
424
425 for span in &trace.spans {
427 let span_id = hex_from_str(
428 &format!("{}-{}", trace.request_id, span.stage.name()),
429 16,
430 );
431 let end_nanos = span.start_unix_nanos
432 + span.duration.as_nanos();
433 let error_attr = match &span.error {
434 Some(e) => format!(
435 r#",{{"key":"error","value":{{"stringValue":"{e}"}}}}"#,
436 e = e.replace('"', "'")
437 ),
438 None => String::new(),
439 };
440 span_jsons.push(format!(
441 r#"{{"traceId":"{trace_id}","spanId":"{sid}","parentSpanId":"{root_span_id}","name":"{stage}","kind":1,"startTimeUnixNano":"{start}","endTimeUnixNano":"{end}","status":{{"code":{code}}},"attributes":[{{"key":"stage","value":{{"stringValue":"{stage}"}}}}{error_attr}]}}"#,
442 trace_id = trace_id,
443 sid = span_id,
444 root_span_id = root_span_id,
445 stage = span.stage.name(),
446 start = span.start_unix_nanos,
447 end = end_nanos,
448 code = if span.success { 1 } else { 2 },
449 error_attr = error_attr,
450 ));
451 }
452
453 let spans_array = span_jsons.join(",");
454 let service_name = &self.config.service_name;
455 let payload = format!(
456 r#"{{"resourceSpans":[{{"resource":{{"attributes":[{{"key":"service.name","value":{{"stringValue":"{service_name}"}}}}]}},"scopeSpans":[{{"spans":[{spans_array}]}}]}}]}}"#
457 );
458 Ok(payload)
459 }
460
461 #[cfg(feature = "zipkin")]
464 async fn export_zipkin(&self, trace: &RequestTrace) -> Result<(), TraceExportError> {
465 let endpoint = self
466 .config
467 .resolved_endpoint()
468 .ok_or_else(|| TraceExportError::NotConfigured("ZIPKIN_ENDPOINT".to_string()))?;
469
470 let payload = self.build_zipkin_payload(trace)?;
471 let url = format!("{endpoint}/api/v2/spans");
472
473 let resp = self
474 .http_client
475 .post(&url)
476 .header("Content-Type", "application/json")
477 .body(payload)
478 .send()
479 .await
480 .map_err(|e| TraceExportError::Http(e.to_string()))?;
481
482 if !resp.status().is_success() {
483 warn!(
484 status = %resp.status(),
485 "Zipkin trace export received non-2xx response"
486 );
487 }
488 Ok(())
489 }
490
491 #[cfg(feature = "zipkin")]
492 fn build_zipkin_payload(&self, trace: &RequestTrace) -> Result<String, TraceExportError> {
493 let trace_id = hex_from_str(&trace.request_id, 32);
494 let root_id = hex_from_str(&trace.request_id, 16);
495 let service_name = &self.config.service_name;
496
497 let root_start_us = trace.entry_unix_nanos / 1_000;
499 let root_duration_us = trace.total_duration.as_micros();
500
501 let mut spans: Vec<String> = vec![format!(
502 r#"{{"traceId":"{trace_id}","id":"{root_id}","name":"request","timestamp":{root_start_us},"duration":{root_duration_us},"localEndpoint":{{"serviceName":"{service_name}"}}}}"#
503 )];
504
505 for span in &trace.spans {
506 let span_id = hex_from_str(
507 &format!("{}-{}", trace.request_id, span.stage.name()),
508 16,
509 );
510 let start_us = span.start_unix_nanos / 1_000;
511 let dur_us = span.duration.as_micros();
512 spans.push(format!(
513 r#"{{"traceId":"{trace_id}","id":"{span_id}","parentId":"{root_id}","name":"{stage}","timestamp":{start_us},"duration":{dur_us},"localEndpoint":{{"serviceName":"{service_name}"}}}}"#,
514 stage = span.stage.name()
515 ));
516 }
517
518 Ok(format!("[{}]", spans.join(",")))
519 }
520
521 #[cfg(feature = "datadog")]
524 async fn export_datadog(&self, trace: &RequestTrace) -> Result<(), TraceExportError> {
525 let endpoint = self
526 .config
527 .resolved_endpoint()
528 .ok_or_else(|| TraceExportError::NotConfigured("DD_AGENT_HOST".to_string()))?;
529
530 let payload = self.build_datadog_payload(trace)?;
531 let url = format!("{endpoint}/v0.4/traces");
532
533 let resp = self
534 .http_client
535 .put(&url)
536 .header("Content-Type", "application/json")
537 .header("X-Datadog-Trace-Count", "1")
538 .body(payload)
539 .send()
540 .await
541 .map_err(|e| TraceExportError::Http(e.to_string()))?;
542
543 if !resp.status().is_success() {
544 warn!(
545 status = %resp.status(),
546 "DataDog trace export received non-2xx response"
547 );
548 }
549 Ok(())
550 }
551
552 #[cfg(feature = "datadog")]
553 fn build_datadog_payload(&self, trace: &RequestTrace) -> Result<String, TraceExportError> {
554 let trace_id: u64 = trace
556 .request_id
557 .bytes()
558 .fold(0u64, |acc, b| acc.wrapping_mul(31).wrapping_add(b as u64));
559 let root_span_id: u64 = trace_id.wrapping_add(1);
560 let service = &self.config.service_name;
561
562 let mut dd_spans: Vec<String> = vec![format!(
563 r#"{{"service":"{service}","name":"request","resource":"{rid}","trace_id":{trace_id},"span_id":{root_span_id},"parent_id":0,"start":{start},"duration":{dur},"error":{err}}}"#,
564 rid = trace.request_id,
565 start = trace.entry_unix_nanos,
566 dur = trace.total_duration.as_nanos(),
567 err = if trace.success { 0 } else { 1 },
568 )];
569
570 for span in &trace.spans {
571 let span_id: u64 = trace_id
572 .wrapping_add(span.stage.index() as u64 + 2);
573 dd_spans.push(format!(
574 r#"{{"service":"{service}","name":"{stage}","resource":"{stage}","trace_id":{trace_id},"span_id":{span_id},"parent_id":{root_span_id},"start":{start},"duration":{dur},"error":{err}}}"#,
575 stage = span.stage.name(),
576 start = span.start_unix_nanos,
577 dur = span.duration.as_nanos(),
578 err = if span.success { 0 } else { 1 },
579 ));
580 }
581
582 Ok(format!("[[{}]]", dd_spans.join(",")))
583 }
584}
585
586fn hex_from_str(s: &str, len: usize) -> String {
591 let mut hash: u128 = 0xcafe_babe_dead_beef_0123_4567_89ab_cdef;
592 for b in s.bytes() {
593 hash = hash.wrapping_mul(6_364_136_223_846_793_005).wrapping_add(b as u128);
594 }
595 format!("{hash:032x}").chars().take(len).collect()
596}
597
598#[cfg(test)]
601mod tests {
602 use super::*;
603
604 #[test]
605 fn pipeline_span_finish_records_duration() {
606 let mut span = PipelineSpan::start("req-1", "sess-1", PipelineStage::Inference);
607 std::thread::sleep(Duration::from_millis(5));
608 span.finish(true, None);
609 assert!(span.duration.as_millis() >= 5);
610 assert!(span.success);
611 }
612
613 #[test]
614 fn request_trace_add_spans() {
615 let mut trace = RequestTrace::new("req-1", "sess-1", "test-svc");
616 let mut span = PipelineSpan::start("req-1", "sess-1", PipelineStage::Rag);
617 span.finish(true, None);
618 trace.add_span(span);
619 trace.finish(true);
620 assert_eq!(trace.spans.len(), 1);
621 assert!(trace.success);
622 }
623
624 #[test]
625 fn otlp_payload_is_valid_json() {
626 let config = TraceExporterConfig {
627 backend: TraceBackend::NoOp,
628 service_name: "test".to_string(),
629 endpoint_override: Some("http://localhost:4318".to_string()),
630 export_timeout: Duration::from_secs(5),
631 };
632 let exporter = TraceExporter::new(config).expect("client should build");
633 let mut trace = RequestTrace::new("req-abc", "sess-1", "test");
634 let mut span = PipelineSpan::start("req-abc", "sess-1", PipelineStage::Inference);
635 span.finish(true, None);
636 trace.add_span(span);
637 trace.finish(true);
638
639 let payload = exporter.build_otlp_payload(&trace).expect("payload ok");
641 let parsed: serde_json::Value =
642 serde_json::from_str(&payload).expect("payload should be valid JSON");
643 assert!(parsed.get("resourceSpans").is_some());
644 }
645
646 #[test]
647 fn hex_from_str_stable() {
648 let a = hex_from_str("hello", 32);
649 let b = hex_from_str("hello", 32);
650 assert_eq!(a, b);
651 assert_eq!(a.len(), 32);
652 }
653
654 #[tokio::test]
655 async fn noop_export_succeeds() {
656 let config = TraceExporterConfig {
657 backend: TraceBackend::NoOp,
658 ..Default::default()
659 };
660 let exporter = TraceExporter::new(config).expect("client ok");
661 let trace = RequestTrace::new("req-noop", "sess-1", "test");
662 exporter.export(&trace).await.expect("noop export ok");
663 }
664}