use std::sync::Arc;
use opentelemetry::Context;
use opentelemetry::trace::Span as _;
use opentelemetry_sdk::Resource;
use opentelemetry_sdk::error::OTelSdkResult;
use opentelemetry_sdk::trace::{Span, SpanData, SpanProcessor};
use super::connection::SharedEngineConnection;
use super::span_exporter::serialize_span_start_snapshot;
use super::types::PREFIX_TRACES;
pub struct LiveSpanStartProcessor {
connection: Arc<SharedEngineConnection>,
resource: Resource,
}
impl LiveSpanStartProcessor {
pub fn new(connection: Arc<SharedEngineConnection>, resource: Resource) -> Self {
Self {
connection,
resource,
}
}
}
impl std::fmt::Debug for LiveSpanStartProcessor {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("LiveSpanStartProcessor").finish()
}
}
impl SpanProcessor for LiveSpanStartProcessor {
fn on_start(&self, span: &mut Span, _cx: &Context) {
if !span.span_context().is_sampled() {
return;
}
let Some(data) = span.exported_data() else {
return;
};
let json = serialize_span_start_snapshot(&data, &self.resource);
let _ = self.connection.send(PREFIX_TRACES, json);
}
fn on_end(&self, _span: SpanData) {}
fn force_flush(&self) -> OTelSdkResult {
Ok(())
}
fn shutdown_with_timeout(&self, _timeout: std::time::Duration) -> OTelSdkResult {
Ok(())
}
}
#[cfg(test)]
mod tests {
use std::sync::{Arc, Mutex};
use opentelemetry::baggage::BaggageExt;
use opentelemetry::trace::{Span as _, Tracer, TracerProvider};
use opentelemetry::{Context, KeyValue};
use opentelemetry_sdk::Resource;
use opentelemetry_sdk::error::OTelSdkResult;
use opentelemetry_sdk::trace::{
InMemorySpanExporter, SdkTracerProvider, SimpleSpanProcessor, Span, SpanData, SpanProcessor,
};
use super::super::baggage_span_processor::BaggageSpanProcessor;
use super::super::span_exporter::serialize_span_start_snapshot;
fn span_data_with_baggage() -> opentelemetry_sdk::trace::SpanData {
let exporter = InMemorySpanExporter::default();
let provider = SdkTracerProvider::builder()
.with_span_processor(BaggageSpanProcessor)
.with_span_processor(SimpleSpanProcessor::new(exporter.clone()))
.build();
let tracer = provider.tracer("test");
let cx = Context::new().with_baggage(vec![
KeyValue::new("iii.session.id", "S-1"),
KeyValue::new("iii.tag.kind", "harness.turn"),
]);
let span = tracer
.span_builder("harness::turn step")
.start_with_context(&tracer, &cx);
drop(span);
exporter
.get_finished_spans()
.expect("exporter ok")
.into_iter()
.next()
.expect("one span")
}
#[test]
fn start_snapshot_has_zero_end_and_keeps_identity_and_baggage() {
let data = span_data_with_baggage();
let resource = Resource::builder_empty()
.with_attributes([KeyValue::new("service.name", "harness")])
.build();
let bytes = serialize_span_start_snapshot(&data, &resource);
let json: serde_json::Value = serde_json::from_slice(&bytes).expect("valid JSON");
let span = &json["resourceSpans"][0]["scopeSpans"][0]["spans"][0];
assert_eq!(
span["endTimeUnixNano"], "0",
"zero end is the pending sentinel"
);
assert_ne!(span["startTimeUnixNano"], "0");
assert_eq!(span["name"], "harness::turn step");
assert_eq!(span["traceId"].as_str().unwrap().len(), 32);
assert_eq!(span["spanId"].as_str().unwrap().len(), 16);
let attrs = span["attributes"].as_array().expect("attributes array");
let has = |key: &str, value: &str| {
attrs
.iter()
.any(|a| a["key"] == key && a["value"]["stringValue"] == value)
};
assert!(has("iii.session.id", "S-1"), "attrs: {attrs:?}");
assert!(has("iii.tag.kind", "harness.turn"), "attrs: {attrs:?}");
let res_attrs = json["resourceSpans"][0]["resource"]["attributes"]
.as_array()
.expect("resource attributes");
assert!(
res_attrs
.iter()
.any(|a| a["key"] == "service.name" && a["value"]["stringValue"] == "harness")
);
}
#[derive(Debug)]
struct CaptureStartFrames(Arc<Mutex<Vec<Vec<u8>>>>, Resource);
impl SpanProcessor for CaptureStartFrames {
fn on_start(&self, span: &mut Span, _cx: &Context) {
if !span.span_context().is_sampled() {
return;
}
let Some(data) = span.exported_data() else {
return;
};
self.0
.lock()
.unwrap()
.push(serialize_span_start_snapshot(&data, &self.1));
}
fn on_end(&self, _span: SpanData) {}
fn force_flush(&self) -> OTelSdkResult {
Ok(())
}
fn shutdown_with_timeout(&self, _t: std::time::Duration) -> OTelSdkResult {
Ok(())
}
}
#[test]
fn on_start_snapshot_sees_baggage_stamped_by_earlier_processor() {
let frames = Arc::new(Mutex::new(Vec::new()));
let resource = Resource::builder_empty().build();
let provider = SdkTracerProvider::builder()
.with_span_processor(BaggageSpanProcessor)
.with_span_processor(CaptureStartFrames(frames.clone(), resource))
.build();
let tracer = provider.tracer("test");
let cx = Context::new().with_baggage(vec![
KeyValue::new("iii.session.id", "S-live"),
KeyValue::new("iii.tag.kind", "harness.turn"),
]);
let span = tracer
.span_builder("harness::turn step")
.start_with_context(&tracer, &cx);
let captured = frames.lock().unwrap().clone();
drop(span);
assert_eq!(captured.len(), 1, "one snapshot per span start");
let json: serde_json::Value = serde_json::from_slice(&captured[0]).unwrap();
let snap = &json["resourceSpans"][0]["scopeSpans"][0]["spans"][0];
assert_eq!(snap["endTimeUnixNano"], "0");
let attrs = snap["attributes"].as_array().unwrap();
assert!(
attrs
.iter()
.any(|a| a["key"] == "iii.session.id" && a["value"]["stringValue"] == "S-live"),
"baggage must be on the start snapshot; attrs: {attrs:?}"
);
}
}