Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
13 changes: 13 additions & 0 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

7 changes: 4 additions & 3 deletions Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -40,9 +40,10 @@ futures-util = "0.3"
tracing = "0.1"
tracing-subscriber = { version = "0.3", features = ["env-filter", "json"] }
tracing-opentelemetry = "0.27"
opentelemetry = { version = "0.26", features = ["trace"] }
opentelemetry_sdk = { version = "0.26", features = ["rt-tokio"] }
opentelemetry-otlp = { version = "0.26", features = ["trace", "grpc-tonic"] }
opentelemetry = { version = "0.26", features = ["trace", "metrics", "logs"] }
opentelemetry_sdk = { version = "0.26", features = ["rt-tokio", "metrics", "logs"] }
opentelemetry-otlp = { version = "0.26", features = ["trace", "metrics", "logs", "grpc-tonic"] }
opentelemetry-appender-tracing = "0.26"
jsonwebtoken = { version = "10", features = ["rust_crypto"] }
subtle = "2"
chrono = { version = "0.4", default-features = false, features = ["clock"] }
Expand Down
138 changes: 112 additions & 26 deletions src/bin/arbitus.rs
Original file line number Diff line number Diff line change
Expand Up @@ -37,7 +37,7 @@ use regex::Regex;
use rusqlite::{Connection, types::Value};
use std::{collections::HashMap, sync::Arc, time::Duration};
use tokio::sync::watch;
use tracing_subscriber::{EnvFilter, layer::SubscriberExt, util::SubscriberInitExt};
use tracing_subscriber::{EnvFilter, Layer, layer::SubscriberExt, util::SubscriberInitExt};

// ── CLI definition ─────────────────────────────────────────────────────────────

Expand Down Expand Up @@ -1198,11 +1198,20 @@ audits: []

// ── OpenTelemetry ──────────────────────────────────────────────────────────────

struct OtelGuard;
struct OtelGuard {
metrics_provider: Option<opentelemetry_sdk::metrics::SdkMeterProvider>,
logger_provider: Option<opentelemetry_sdk::logs::LoggerProvider>,
}

impl Drop for OtelGuard {
fn drop(&mut self) {
opentelemetry::global::shutdown_tracer_provider();
if let Some(p) = self.metrics_provider.take() {
let _ = p.shutdown();
}
if let Some(p) = self.logger_provider.take() {
let _ = p.shutdown();
}
}
}

Expand All @@ -1213,35 +1222,64 @@ fn init_tracing(telemetry: Option<&TelemetryConfig>) -> Option<OtelGuard> {
let tracer = telemetry.and_then(|tel| match build_otel_tracer(tel) {
Ok(t) => Some(t),
Err(e) => {
eprintln!("warn: OTel init failed: {e}");
eprintln!("warn: OTel traces init failed: {e}");
None
}
});

let has_otel = tracer.is_some();

match (json, tracer) {
(true, Some(t)) => tracing_subscriber::registry()
.with(filter)
.with(tracing_subscriber::fmt::layer().json())
.with(tracing_opentelemetry::layer().with_tracer(t))
.init(),
(true, None) => tracing_subscriber::registry()
.with(filter)
.with(tracing_subscriber::fmt::layer().json())
.init(),
(false, Some(t)) => tracing_subscriber::registry()
.with(filter)
.with(tracing_subscriber::fmt::layer())
.with(tracing_opentelemetry::layer().with_tracer(t))
.init(),
(false, None) => tracing_subscriber::registry()
.with(filter)
.with(tracing_subscriber::fmt::layer())
.init(),
}
let metrics_provider =
telemetry
.filter(|t| t.export_metrics)
.and_then(|tel| match build_otel_metrics(tel) {
Ok(p) => {
opentelemetry::global::set_meter_provider(p.clone());
Some(p)
}
Err(e) => {
eprintln!("warn: OTel metrics init failed: {e}");
None
}
});

let logger_provider =
telemetry
.filter(|t| t.export_logs)
.and_then(|tel| match build_otel_logs(tel) {
Ok(p) => Some(p),
Err(e) => {
eprintln!("warn: OTel logs init failed: {e}");
None
}
});

let has_traces = tracer.is_some();
let otel_trace_layer = tracer.map(|t| tracing_opentelemetry::layer().with_tracer(t));
let otel_log_layer = logger_provider
.as_ref()
.map(opentelemetry_appender_tracing::layer::OpenTelemetryTracingBridge::new);

let fmt_layer = if json {
tracing_subscriber::fmt::layer().json().boxed()
} else {
tracing_subscriber::fmt::layer().boxed()
};

if has_otel { Some(OtelGuard) } else { None }
tracing_subscriber::registry()
.with(filter)
.with(fmt_layer)
.with(otel_trace_layer)
.with(otel_log_layer)
.init();

let any_otel = has_traces || metrics_provider.is_some() || logger_provider.is_some();
if any_otel {
Some(OtelGuard {
metrics_provider,
logger_provider,
})
} else {
None
}
}

fn build_otel_tracer(tel: &TelemetryConfig) -> anyhow::Result<opentelemetry_sdk::trace::Tracer> {
Expand All @@ -1268,3 +1306,51 @@ fn build_otel_tracer(tel: &TelemetryConfig) -> anyhow::Result<opentelemetry_sdk:

Ok(provider.tracer("arbitus"))
}

fn build_otel_metrics(
tel: &TelemetryConfig,
) -> anyhow::Result<opentelemetry_sdk::metrics::SdkMeterProvider> {
use opentelemetry::KeyValue;
use opentelemetry_otlp::WithExportConfig;
use opentelemetry_sdk::Resource;

let resource = Resource::new(vec![KeyValue::new(
"service.name",
tel.service_name.clone(),
)]);

opentelemetry_otlp::new_pipeline()
.metrics(opentelemetry_sdk::runtime::Tokio)
.with_exporter(
opentelemetry_otlp::new_exporter()
.tonic()
.with_endpoint(&tel.otlp_endpoint),
)
.with_resource(resource)
.build()
.map_err(|e| anyhow::anyhow!("OTLP metrics pipeline: {e}"))
}

fn build_otel_logs(
tel: &TelemetryConfig,
) -> anyhow::Result<opentelemetry_sdk::logs::LoggerProvider> {
use opentelemetry::KeyValue;
use opentelemetry_otlp::WithExportConfig;
use opentelemetry_sdk::Resource;

let resource = Resource::new(vec![KeyValue::new(
"service.name",
tel.service_name.clone(),
)]);

opentelemetry_otlp::new_pipeline()
.logging()
.with_resource(resource)
.with_exporter(
opentelemetry_otlp::new_exporter()
.tonic()
.with_endpoint(&tel.otlp_endpoint),
)
.install_batch(opentelemetry_sdk::runtime::Tokio)
.map_err(|e| anyhow::anyhow!("OTLP logs pipeline: {e}"))
}
49 changes: 47 additions & 2 deletions src/config.rs
Original file line number Diff line number Diff line change
Expand Up @@ -538,15 +538,21 @@ fn default_agent_claim() -> String {

// ── Telemetry ─────────────────────────────────────────────────────────────────

/// OpenTelemetry tracing configuration.
/// When set, spans are exported to the configured OTLP endpoint.
/// OpenTelemetry configuration.
/// When set, spans, metrics, and/or logs are exported to the configured OTLP endpoint.
#[derive(Debug, Deserialize, Clone)]
pub struct TelemetryConfig {
/// OTLP gRPC endpoint (e.g. `http://localhost:4317`).
pub otlp_endpoint: String,
/// `service.name` resource attribute. Defaults to `"arbitus"`.
#[serde(default = "default_service_name")]
pub service_name: String,
/// Export metrics via OTLP alongside the Prometheus `/metrics` endpoint.
#[serde(default)]
pub export_metrics: bool,
/// Export structured logs via OTLP (bridges all `tracing` events).
#[serde(default)]
pub export_logs: bool,
}

fn default_service_name() -> String {
Expand Down Expand Up @@ -1040,4 +1046,43 @@ mod tests {
);
assert!(cfg.validate().is_ok());
}

// ── TelemetryConfig ───────────────────────────────────────────────────────

#[test]
fn telemetry_config_defaults() {
let yaml = r#"
otlp_endpoint: "http://localhost:4317"
"#;
let cfg: TelemetryConfig = serde_yaml::from_str(yaml).unwrap();
assert_eq!(cfg.service_name, "arbitus");
assert!(!cfg.export_metrics);
assert!(!cfg.export_logs);
}

#[test]
fn telemetry_config_all_fields() {
let yaml = r#"
otlp_endpoint: "http://otel-collector:4317"
service_name: "my-service"
export_metrics: true
export_logs: true
"#;
let cfg: TelemetryConfig = serde_yaml::from_str(yaml).unwrap();
assert_eq!(cfg.otlp_endpoint, "http://otel-collector:4317");
assert_eq!(cfg.service_name, "my-service");
assert!(cfg.export_metrics);
assert!(cfg.export_logs);
}

#[test]
fn telemetry_config_partial_flags() {
let yaml = r#"
otlp_endpoint: "http://localhost:4317"
export_metrics: true
"#;
let cfg: TelemetryConfig = serde_yaml::from_str(yaml).unwrap();
assert!(cfg.export_metrics);
assert!(!cfg.export_logs);
}
}
68 changes: 68 additions & 0 deletions src/metrics.rs
Original file line number Diff line number Diff line change
@@ -1,3 +1,4 @@
use opentelemetry::{KeyValue, global};
use prometheus::{Counter, CounterVec, Encoder, Opts, Registry, TextEncoder};

pub struct GatewayMetrics {
Expand Down Expand Up @@ -60,6 +61,18 @@ impl GatewayMetrics {

pub fn record(&self, agent: &str, outcome: &str) {
self.requests.with_label_values(&[agent, outcome]).inc();
// Mirror to OTLP when a global meter provider is installed (no-op otherwise).
global::meter("arbitus")
.u64_counter("arbitus.requests.total")
.with_description("Total requests processed by arbitus")
.init()
.add(
1,
&[
KeyValue::new("agent", agent.to_string()),
KeyValue::new("outcome", outcome.to_string()),
],
);
}

/// Record estimated token usage for a single request.
Expand All @@ -71,11 +84,33 @@ impl GatewayMetrics {
self.tokens
.with_label_values(&[agent, "input"])
.inc_by(f64::from(input_tokens));
global::meter("arbitus")
.f64_counter("arbitus.tokens.total")
.with_description("Estimated tokens processed by arbitus")
.init()
.add(
f64::from(input_tokens),
&[
KeyValue::new("agent", agent.to_string()),
KeyValue::new("direction", "input"),
],
);
}
if output_tokens > 0 {
self.tokens
.with_label_values(&[agent, "output"])
.inc_by(f64::from(output_tokens));
global::meter("arbitus")
.f64_counter("arbitus.tokens.total")
.with_description("Estimated tokens processed by arbitus")
.init()
.add(
f64::from(output_tokens),
&[
KeyValue::new("agent", agent.to_string()),
KeyValue::new("direction", "output"),
],
);
}
}

Expand Down Expand Up @@ -133,4 +168,37 @@ mod tests {
assert!(rendered.contains(r#"agent="cursor""#));
assert!(rendered.contains(r#"agent="claude""#));
}

// ── OTel dual-export (no-op without provider) ─────────────────────────────

#[test]
fn record_does_not_panic_without_otel_provider() {
// No global OTel meter provider installed — calls must be silent no-ops.
let m = GatewayMetrics::new().unwrap();
m.record("cursor", "allowed");
m.record("cursor", "blocked");
// Prometheus counter still incremented
let rendered = m.render();
assert!(rendered.contains("arbitus_requests_total"));
}

#[test]
fn record_tokens_does_not_panic_without_otel_provider() {
let m = GatewayMetrics::new().unwrap();
m.record_tokens("cursor", 100, 200);
let rendered = m.render();
assert!(rendered.contains("arbitus_tokens_total"));
}

#[test]
fn prometheus_metrics_unaffected_by_otel_calls() {
let m = GatewayMetrics::new().unwrap();
m.record("agent-a", "allowed");
m.record("agent-a", "allowed");
m.record("agent-a", "blocked");
let rendered = m.render();
// Two allowed + one blocked — Prometheus counters must reflect this.
assert!(rendered.contains(r#"outcome="allowed""#));
assert!(rendered.contains(r#"outcome="blocked""#));
}
}
Loading
Loading