use std::time::Duration;
use async_trait::async_trait;
use opentelemetry::logs::{AnyValue, LogRecord, Logger, LoggerProvider, Severity};
use opentelemetry::KeyValue;
use opentelemetry_otlp::{LogExporter, WithExportConfig, WithHttpConfig, WithTonicConfig};
use opentelemetry_sdk::logs::SdkLoggerProvider;
use opentelemetry_sdk::Resource;
use secrecy::{ExposeSecret, SecretString};
use crate::error::TelemetryError;
use crate::event::{severity_of, Event};
use crate::sink::TelemetrySink;
#[derive(Debug, Clone)]
pub struct OtlpSinkConfig {
pub endpoint: String,
pub headers: Vec<(String, SecretString)>,
pub timeout: Duration,
pub resource_attrs: Vec<(String, String)>,
}
impl Default for OtlpSinkConfig {
fn default() -> Self {
Self {
endpoint: "http://127.0.0.1:4317".into(),
headers: Vec::new(),
timeout: Duration::from_secs(10),
resource_attrs: Vec::new(),
}
}
}
#[derive(Debug)]
pub struct OtlpSink {
provider: SdkLoggerProvider,
logger_name: &'static str,
}
impl OtlpSink {
pub fn new(config: OtlpSinkConfig) -> Result<Self, TelemetryError> {
let transport = Transport::from_endpoint(&config.endpoint)?;
let exporter = build_exporter(&config, &transport)?;
let OtlpSinkConfig { resource_attrs, .. } = config;
let mut resource = Resource::builder().with_service_name("rtb-telemetry");
for (k, v) in resource_attrs {
resource = resource.with_attribute(KeyValue::new(k, v));
}
let provider = SdkLoggerProvider::builder()
.with_resource(resource.build())
.with_batch_exporter(exporter)
.build();
Ok(Self { provider, logger_name: "rtb-telemetry" })
}
}
#[async_trait]
impl TelemetrySink for OtlpSink {
async fn emit(&self, event: &Event) -> Result<(), TelemetryError> {
let redacted = event.redacted();
let body =
serde_json::to_string(&redacted).map_err(|e| TelemetryError::Serde(e.to_string()))?;
let severity = match severity_of(&redacted) {
"ERROR" => Severity::Error,
_ => Severity::Info,
};
let logger = self.provider.logger(self.logger_name);
let mut record = logger.create_log_record();
record.set_event_name("rtb.telemetry.event");
record.set_severity_number(severity);
record.set_severity_text(severity_of(&redacted));
record.set_body(AnyValue::String(body.into()));
record.add_attribute("tool", redacted.tool.clone());
record.add_attribute("tool.version", redacted.tool_version.clone());
record.add_attribute("event.name", redacted.name.clone());
logger.emit(record);
Ok(())
}
async fn flush(&self) -> Result<(), TelemetryError> {
self.provider.force_flush().map_err(|e| TelemetryError::Otlp(format!("flush: {e}")))
}
}
enum Transport {
Grpc,
Http,
}
impl Transport {
fn from_endpoint(endpoint: &str) -> Result<Self, TelemetryError> {
if let Some(rest) =
endpoint.strip_prefix("grpc://").or_else(|| endpoint.strip_prefix("grpcs://"))
{
if rest.is_empty() {
return Err(TelemetryError::Otlp(format!("empty host in endpoint {endpoint:?}")));
}
return Ok(Self::Grpc);
}
if endpoint.starts_with("http://") || endpoint.starts_with("https://") {
let is_grpc_port = endpoint.contains(":4317");
return Ok(if is_grpc_port { Self::Grpc } else { Self::Http });
}
Err(TelemetryError::Otlp(format!(
"unsupported endpoint scheme in {endpoint:?} \
(expected grpc://, grpcs://, http://, or https://)"
)))
}
}
fn build_exporter(
config: &OtlpSinkConfig,
transport: &Transport,
) -> Result<LogExporter, TelemetryError> {
let endpoint = config.endpoint.replace("grpcs://", "https://").replace("grpc://", "http://");
match transport {
Transport::Grpc => {
let mut builder = LogExporter::builder()
.with_tonic()
.with_endpoint(endpoint)
.with_timeout(config.timeout);
if !config.headers.is_empty() {
let mut metadata = tonic::metadata::MetadataMap::new();
for (k, v) in &config.headers {
let key =
tonic::metadata::MetadataKey::from_bytes(k.as_bytes()).map_err(|e| {
TelemetryError::Otlp(format!("invalid header name {k:?}: {e}"))
})?;
let val = v
.expose_secret()
.parse()
.map_err(|e| TelemetryError::Otlp(format!("invalid header value: {e}")))?;
metadata.insert(key, val);
}
builder = builder.with_metadata(metadata);
}
builder.build().map_err(|e| TelemetryError::Otlp(format!("build: {e}")))
}
Transport::Http => {
let mut builder = LogExporter::builder()
.with_http()
.with_endpoint(endpoint)
.with_timeout(config.timeout);
if !config.headers.is_empty() {
let map: std::collections::HashMap<String, String> = config
.headers
.iter()
.map(|(k, v)| (k.clone(), v.expose_secret().to_string()))
.collect();
builder = builder.with_headers(map);
}
builder.build().map_err(|e| TelemetryError::Otlp(format!("build: {e}")))
}
}
}