use std::time::{Duration, Instant};
use tracing::{field, info_span, Span};
#[derive(Debug, Clone)]
pub struct TracingConfig {
pub service_name: String,
pub record_payloads: bool,
pub record_responses: bool,
pub max_payload_size: usize,
}
impl Default for TracingConfig {
fn default() -> Self {
Self {
service_name: "talos-client".to_string(),
record_payloads: false,
record_responses: false,
max_payload_size: 4096,
}
}
}
impl TracingConfig {
pub fn builder() -> TracingConfigBuilder {
TracingConfigBuilder::default()
}
}
#[derive(Debug, Default)]
pub struct TracingConfigBuilder {
service_name: Option<String>,
record_payloads: Option<bool>,
record_responses: Option<bool>,
max_payload_size: Option<usize>,
}
impl TracingConfigBuilder {
pub fn service_name(mut self, name: impl Into<String>) -> Self {
self.service_name = Some(name.into());
self
}
pub fn record_payloads(mut self, enabled: bool) -> Self {
self.record_payloads = Some(enabled);
self
}
pub fn record_responses(mut self, enabled: bool) -> Self {
self.record_responses = Some(enabled);
self
}
pub fn max_payload_size(mut self, size: usize) -> Self {
self.max_payload_size = Some(size);
self
}
pub fn build(self) -> TracingConfig {
let default = TracingConfig::default();
TracingConfig {
service_name: self.service_name.unwrap_or(default.service_name),
record_payloads: self.record_payloads.unwrap_or(default.record_payloads),
record_responses: self.record_responses.unwrap_or(default.record_responses),
max_payload_size: self.max_payload_size.unwrap_or(default.max_payload_size),
}
}
}
#[derive(Debug)]
pub struct TalosSpan {
span: Span,
start: Instant,
method: String,
endpoint: String,
}
impl TalosSpan {
pub fn new(method: &str, endpoint: &str) -> Self {
let span = info_span!(
"talos.grpc",
rpc.system = "grpc",
rpc.service = "talos.machine.MachineService",
rpc.method = %method,
server.address = %endpoint,
rpc.grpc.status_code = field::Empty,
otel.status_code = field::Empty,
error.message = field::Empty,
duration_ms = field::Empty,
);
Self {
span,
start: Instant::now(),
method: method.to_string(),
endpoint: endpoint.to_string(),
}
}
pub fn with_service(method: &str, service: &str, endpoint: &str) -> Self {
let span = info_span!(
"talos.grpc",
rpc.system = "grpc",
rpc.service = %service,
rpc.method = %method,
server.address = %endpoint,
rpc.grpc.status_code = field::Empty,
otel.status_code = field::Empty,
error.message = field::Empty,
duration_ms = field::Empty,
);
Self {
span,
start: Instant::now(),
method: method.to_string(),
endpoint: endpoint.to_string(),
}
}
pub fn span(&self) -> &Span {
&self.span
}
pub fn method(&self) -> &str {
&self.method
}
pub fn endpoint(&self) -> &str {
&self.endpoint
}
pub fn elapsed(&self) -> Duration {
self.start.elapsed()
}
pub fn record_success(&self, duration: Duration) {
self.span.record("rpc.grpc.status_code", 0i64); self.span.record("otel.status_code", "OK");
self.span.record("duration_ms", duration.as_millis() as i64);
}
pub fn record_error(&self, error: &str) {
let duration = self.start.elapsed();
self.span.record("rpc.grpc.status_code", 2i64); self.span.record("otel.status_code", "ERROR");
self.span.record("error.message", error);
self.span.record("duration_ms", duration.as_millis() as i64);
}
pub fn record_grpc_status(&self, code: i32) {
self.span.record("rpc.grpc.status_code", code as i64);
let status = if code == 0 { "OK" } else { "ERROR" };
self.span.record("otel.status_code", status);
self.span
.record("duration_ms", self.start.elapsed().as_millis() as i64);
}
pub fn enter(&self) -> tracing::span::Entered<'_> {
self.span.enter()
}
}
#[macro_export]
macro_rules! instrument_talos {
($method:expr, $endpoint:expr, $body:expr) => {{
let span = $crate::runtime::tracing::TalosSpan::new($method, $endpoint);
let _guard = span.enter();
let start = std::time::Instant::now();
let result = $body;
let duration = start.elapsed();
match &result {
Ok(_) => span.record_success(duration),
Err(e) => span.record_error(&format!("{}", e)),
}
result
}};
}
#[derive(Debug, Clone)]
pub struct SpanFactory {
config: TracingConfig,
}
impl SpanFactory {
pub fn new(config: TracingConfig) -> Self {
Self { config }
}
pub fn create_span(&self, method: &str, endpoint: &str) -> TalosSpan {
TalosSpan::with_service(method, "talos.machine.MachineService", endpoint)
}
pub fn create_etcd_span(&self, method: &str, endpoint: &str) -> TalosSpan {
TalosSpan::with_service(method, "talos.machine.MachineService/Etcd", endpoint)
}
pub fn config(&self) -> &TracingConfig {
&self.config
}
}
impl Default for SpanFactory {
fn default() -> Self {
Self::new(TracingConfig::default())
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn test_tracing_config_default() {
let config = TracingConfig::default();
assert_eq!(config.service_name, "talos-client");
assert!(!config.record_payloads);
assert!(!config.record_responses);
assert_eq!(config.max_payload_size, 4096);
}
#[test]
fn test_tracing_config_builder() {
let config = TracingConfig::builder()
.service_name("my-service")
.record_payloads(true)
.record_responses(true)
.max_payload_size(8192)
.build();
assert_eq!(config.service_name, "my-service");
assert!(config.record_payloads);
assert!(config.record_responses);
assert_eq!(config.max_payload_size, 8192);
}
#[test]
fn test_talos_span_new() {
let span = TalosSpan::new("Version", "10.0.0.1:50000");
assert_eq!(span.method(), "Version");
assert_eq!(span.endpoint(), "10.0.0.1:50000");
}
#[test]
fn test_talos_span_with_service() {
let span = TalosSpan::with_service(
"EtcdMemberList",
"talos.machine.MachineService/Etcd",
"10.0.0.1:50000",
);
assert_eq!(span.method(), "EtcdMemberList");
}
#[test]
fn test_talos_span_record_success() {
let span = TalosSpan::new("Version", "10.0.0.1:50000");
span.record_success(Duration::from_millis(42));
}
#[test]
fn test_talos_span_record_error() {
let span = TalosSpan::new("Version", "10.0.0.1:50000");
span.record_error("Connection refused");
}
#[test]
fn test_talos_span_record_grpc_status() {
let span = TalosSpan::new("Version", "10.0.0.1:50000");
span.record_grpc_status(0); span.record_grpc_status(14); }
#[test]
fn test_span_factory_new() {
let config = TracingConfig::builder()
.service_name("test-service")
.build();
let factory = SpanFactory::new(config);
assert_eq!(factory.config().service_name, "test-service");
}
#[test]
fn test_span_factory_create_span() {
let factory = SpanFactory::default();
let span = factory.create_span("Version", "10.0.0.1:50000");
assert_eq!(span.method(), "Version");
}
#[test]
fn test_span_factory_create_etcd_span() {
let factory = SpanFactory::default();
let span = factory.create_etcd_span("EtcdMemberList", "10.0.0.1:50000");
assert_eq!(span.method(), "EtcdMemberList");
}
#[test]
fn test_talos_span_elapsed() {
let span = TalosSpan::new("Version", "10.0.0.1:50000");
std::thread::sleep(Duration::from_millis(10));
let elapsed = span.elapsed();
assert!(elapsed >= Duration::from_millis(10));
}
}