use std::{ sync::{ OnceLock, atomic::{AtomicBool, Ordering}, mpsc::{SyncSender, sync_channel}, }, time::{Duration, SystemTime, UNIX_EPOCH}, }; use serde::Deserialize; use serde_json::json; const DEFAULT_OTLP_LOGS_ENDPOINT: &str = "https://telemetry.cbsk-tech.de/v1/logs"; const DEFAULT_OTLP_TRACES_ENDPOINT: &str = "https://telemetry.cbsk-tech.de/v1/traces"; const DEFAULT_OTLP_METRICS_ENDPOINT: &str = "https://telemetry.cbsk-tech.de/v1/metrics"; const METRICS_INTERVAL: Duration = Duration::from_secs(30); const MAX_BODY_BYTES: usize = 2_048; static ENABLED: AtomicBool = AtomicBool::new(false); static SENDER: OnceLock> = OnceLock::new(); #[derive(Debug)] enum TelemetrySignal { Log(TelemetryRecord), Span(FrontendSpan), Metrics(ProcessMetrics), } #[derive(Debug)] struct ProcessMetrics { cpu_utilization: f64, memory_bytes: u64, } #[derive(Debug)] struct TelemetryRecord { severity: &'static str, severity_number: u8, body: String, event_name: Option, } #[derive(Debug, Deserialize)] #[serde(rename_all = "camelCase")] pub struct FrontendLog { level: String, message: String, event_name: Option, } #[derive(Debug, Deserialize)] #[serde(rename_all = "camelCase")] pub struct FrontendSpan { name: String, trace_id: String, span_id: String, started_at_ms: u64, duration_ms: f64, success: bool, } pub fn init() { let endpoint = std::env::var("GITTY_OTLP_LOGS_ENDPOINT") .unwrap_or_else(|_| DEFAULT_OTLP_LOGS_ENDPOINT.to_string()); let traces_endpoint = std::env::var("GITTY_OTLP_TRACES_ENDPOINT") .unwrap_or_else(|_| DEFAULT_OTLP_TRACES_ENDPOINT.to_string()); let metrics_endpoint = std::env::var("GITTY_OTLP_METRICS_ENDPOINT") .unwrap_or_else(|_| DEFAULT_OTLP_METRICS_ENDPOINT.to_string()); let (sender, receiver) = sync_channel::(256); if SENDER.set(sender.clone()).is_err() { return; } start_process_metrics(sender); std::thread::Builder::new() .name("gitty-telemetry".into()) .spawn(move || { let client = match reqwest::blocking::Client::builder() .timeout(Duration::from_secs(5)) .build() { Ok(client) => client, Err(error) => { eprintln!("[WARN] [gitty::telemetry] could not create OTLP client: {error}"); return; } }; let mut export_warning_shown = false; while let Ok(signal) = receiver.recv() { let (target_endpoint, payload) = match signal { TelemetrySignal::Log(record) => (&endpoint, log_payload(record)), TelemetrySignal::Span(span) => (&traces_endpoint, span_payload(span)), TelemetrySignal::Metrics(metrics) => { (&metrics_endpoint, metrics_payload(metrics)) } }; export( &client, target_endpoint, &payload, &mut export_warning_shown, ); } }) .expect("failed to start telemetry worker"); } fn start_process_metrics(sender: SyncSender) { std::thread::Builder::new() .name("gitty-process-metrics".into()) .spawn(move || { use sysinfo::{ProcessRefreshKind, ProcessesToUpdate, System}; let Ok(pid) = sysinfo::get_current_pid() else { return; }; let refresh_kind = ProcessRefreshKind::nothing().with_cpu().with_memory(); let cpu_count = std::thread::available_parallelism() .map(|count| count.get()) .unwrap_or(1) as f64; let mut system = System::new(); let mut primed = false; loop { std::thread::sleep(METRICS_INTERVAL); if !ENABLED.load(Ordering::Relaxed) { primed = false; continue; } system.refresh_processes_specifics( ProcessesToUpdate::Some(&[pid]), true, refresh_kind, ); let Some(process) = system.process(pid) else { continue; }; if !primed { primed = true; continue; } let metrics = ProcessMetrics { cpu_utilization: (f64::from(process.cpu_usage()) / (100.0 * cpu_count)) .clamp(0.0, 1.0), memory_bytes: process.memory(), }; let _ = sender.try_send(TelemetrySignal::Metrics(metrics)); } }) .expect("failed to start process metrics worker"); } fn resource_attributes() -> serde_json::Value { json!([ { "key": "service.name", "value": { "stringValue": "gitty-desktop" } }, { "key": "service.version", "value": { "stringValue": env!("CARGO_PKG_VERSION") } }, { "key": "deployment.environment.name", "value": { "stringValue": if cfg!(debug_assertions) { "development" } else { "production" } } }, { "key": "os.type", "value": { "stringValue": std::env::consts::OS } }, { "key": "host.arch", "value": { "stringValue": std::env::consts::ARCH } } ]) } fn log_payload(record: TelemetryRecord) -> serde_json::Value { let timestamp = SystemTime::now() .duration_since(UNIX_EPOCH) .unwrap_or_default() .as_nanos() .to_string(); let mut attributes = vec![json!({ "key": "telemetry.sdk.language", "value": { "stringValue": "rust" } })]; if let Some(event_name) = record.event_name { attributes.push(json!({ "key": "event.name", "value": { "stringValue": event_name } })); } let payload = json!({ "resourceLogs": [{ "resource": { "attributes": resource_attributes() }, "scopeLogs": [{ "scope": { "name": "gitty.telemetry" }, "logRecords": [{ "timeUnixNano": timestamp, "observedTimeUnixNano": timestamp, "severityNumber": record.severity_number, "severityText": record.severity, "body": { "stringValue": record.body }, "attributes": attributes }] }] }] }); payload } fn span_payload(mut span: FrontendSpan) -> serde_json::Value { if span.name.len() > 128 { span.name.truncate(128); } let start = u128::from(span.started_at_ms) * 1_000_000; let duration = (span.duration_ms.max(0.0) * 1_000_000.0) as u128; let attributes = vec![ json!({ "key": "rpc.system", "value": { "stringValue": "tauri" } }), json!({ "key": "rpc.method", "value": { "stringValue": span.name } }), ]; json!({ "resourceSpans": [{ "resource": { "attributes": resource_attributes() }, "scopeSpans": [{ "scope": { "name": "gitty.tauri", "version": env!("CARGO_PKG_VERSION") }, "spans": [{ "traceId": span.trace_id, "spanId": span.span_id, "name": format!("tauri.{}", span.name), "kind": 1, "startTimeUnixNano": start.to_string(), "endTimeUnixNano": (start + duration).to_string(), "attributes": attributes, "status": { "code": if span.success { 1 } else { 2 }, "message": if span.success { "" } else { "Tauri command failed" } } }] }] }] }) } fn metrics_payload(metrics: ProcessMetrics) -> serde_json::Value { let timestamp = SystemTime::now() .duration_since(UNIX_EPOCH) .unwrap_or_default() .as_nanos() .to_string(); json!({ "resourceMetrics": [{ "resource": { "attributes": resource_attributes() }, "scopeMetrics": [{ "scope": { "name": "gitty.process", "version": env!("CARGO_PKG_VERSION") }, "metrics": [ { "name": "process.cpu.utilization", "description": "Normalized CPU utilization of the Gitty process", "unit": "1", "gauge": { "dataPoints": [{ "timeUnixNano": timestamp, "asDouble": metrics.cpu_utilization, "attributes": [] }] } }, { "name": "process.memory.usage", "description": "Resident memory used by the Gitty process", "unit": "By", "gauge": { "dataPoints": [{ "timeUnixNano": timestamp, "asInt": metrics.memory_bytes.to_string(), "attributes": [] }] } } ] }] }] }) } fn export( client: &reqwest::blocking::Client, endpoint: &str, payload: &serde_json::Value, export_warning_shown: &mut bool, ) { match client.post(endpoint).json(payload).send() { Ok(response) if response.status().is_success() => { if response .headers() .get(reqwest::header::CONTENT_TYPE) .and_then(|value| value.to_str().ok()) .is_some_and(|value| value.starts_with("text/html")) && !*export_warning_shown { eprintln!( "[WARN] [gitty::telemetry] OTLP endpoint returned HTML; configure the collector endpoint with GITTY_OTLP_LOGS_ENDPOINT" ); *export_warning_shown = true; } } Ok(response) if !*export_warning_shown => { eprintln!( "[WARN] [gitty::telemetry] OTLP export failed with status {}", response.status() ); *export_warning_shown = true; } Err(error) if !*export_warning_shown => { eprintln!("[WARN] [gitty::telemetry] OTLP export failed: {error}"); *export_warning_shown = true; } _ => {} } } fn emit(record: TelemetryRecord) { if !ENABLED.load(Ordering::Relaxed) { return; } if let Some(sender) = SENDER.get() { let _ = sender.try_send(TelemetrySignal::Log(record)); } } #[tauri::command] pub fn set_telemetry_enabled(enabled: bool) { ENABLED.store(enabled, Ordering::Relaxed); if enabled { emit(TelemetryRecord { severity: "INFO", severity_number: 9, body: "Telemetry enabled".into(), event_name: Some("telemetry.enabled".into()), }); } } #[tauri::command] pub fn emit_frontend_log(log: FrontendLog) { let (severity, severity_number) = match log.level.as_str() { "error" => ("ERROR", 17), "warn" => ("WARN", 13), _ => ("INFO", 9), }; let mut body = log.message; if body.len() > MAX_BODY_BYTES { body.truncate(MAX_BODY_BYTES); } emit(TelemetryRecord { severity, severity_number, body, event_name: log.event_name, }); } #[tauri::command] pub fn emit_frontend_span(span: FrontendSpan) { if !ENABLED.load(Ordering::Relaxed) { return; } if span.trace_id.len() != 32 || span.span_id.len() != 16 || !span.trace_id.bytes().all(|byte| byte.is_ascii_hexdigit()) || !span.span_id.bytes().all(|byte| byte.is_ascii_hexdigit()) { return; } if let Some(sender) = SENDER.get() { let _ = sender.try_send(TelemetrySignal::Span(span)); } }