From 23588ab1f4268a85daa0feee7d608c5d5a528798 Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Thu, 19 Mar 2026 14:15:33 -0400 Subject: [PATCH] Add buffered telemetry with track! macro and extract fabro-telemetry crate Extract telemetry from fabro-util into a dedicated fabro-telemetry crate. Replace the synchronous Telemetry struct with a global background buffer that flushes periodically via blocking HTTP (mid-run) and detached subprocess (final flush at exit). The new API is init_cli()/track!()/shutdown(). Co-Authored-By: Claude Opus 4.6 (1M context) --- AGENTS.md | 3 +- Cargo.lock | 35 ++- lib/crates/fabro-cli/Cargo.toml | 1 + lib/crates/fabro-cli/src/main.rs | 73 ++---- lib/crates/fabro-telemetry/Cargo.toml | 31 +++ .../src}/anonymous_id.rs | 0 lib/crates/fabro-telemetry/src/buffer.rs | 221 ++++++++++++++++ .../src}/context.rs | 0 .../src}/event.rs | 0 .../telemetry => fabro-telemetry/src}/git.rs | 0 lib/crates/fabro-telemetry/src/lib.rs | 245 ++++++++++++++++++ .../src}/panic.rs | 29 +-- .../src}/sanitize.rs | 0 lib/crates/fabro-telemetry/src/sender.rs | 67 ++++- .../src}/spawn.rs | 0 lib/crates/fabro-util/Cargo.toml | 10 - lib/crates/fabro-util/src/lib.rs | 1 - lib/crates/fabro-util/src/telemetry/mod.rs | 165 ------------ lib/crates/fabro-util/src/telemetry/sender.rs | 219 ---------------- 19 files changed, 631 insertions(+), 469 deletions(-) create mode 100644 lib/crates/fabro-telemetry/Cargo.toml rename lib/crates/{fabro-util/src/telemetry => fabro-telemetry/src}/anonymous_id.rs (100%) create mode 100644 lib/crates/fabro-telemetry/src/buffer.rs rename lib/crates/{fabro-util/src/telemetry => fabro-telemetry/src}/context.rs (100%) rename lib/crates/{fabro-util/src/telemetry => fabro-telemetry/src}/event.rs (100%) rename lib/crates/{fabro-util/src/telemetry => fabro-telemetry/src}/git.rs (100%) create mode 100644 lib/crates/fabro-telemetry/src/lib.rs rename lib/crates/{fabro-util/src/telemetry => fabro-telemetry/src}/panic.rs (85%) rename lib/crates/{fabro-util/src/telemetry => fabro-telemetry/src}/sanitize.rs (100%) rename lib/crates/{fabro-util/src/telemetry => fabro-telemetry/src}/spawn.rs (100%) delete mode 100644 lib/crates/fabro-util/src/telemetry/mod.rs delete mode 100644 lib/crates/fabro-util/src/telemetry/sender.rs diff --git a/AGENTS.md b/AGENTS.md index 261ce4604..58a847ad7 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -64,7 +64,8 @@ Fabro is an AI-powered workflow orchestration platform. Workflows are defined as - **fabro-slack** — Slack integration (socket mode, blocks API) - **fabro-devcontainer** — Parses `.devcontainer/devcontainer.json` for container setup - **fabro-git-storage** — Git-based storage with branch store and snapshots -- **fabro-util** — Shared utilities (redaction, telemetry, terminal formatting) +- **fabro-telemetry** — CLI analytics (Segment) and crash reporting (Sentry), with anonymous IDs, command sanitization, and detached subprocess delivery +- **fabro-util** — Shared utilities (redaction, terminal formatting) ### TypeScript (`apps/` and `lib/packages/`) - **apps/fabro-web** — React 19 + React Router + Vite + Tailwind CSS frontend diff --git a/Cargo.lock b/Cargo.lock index ce95cc186..6e0fc92eb 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1367,6 +1367,7 @@ dependencies = [ "fabro-openai-oauth", "fabro-retro", "fabro-ssh", + "fabro-telemetry", "fabro-util", "fabro-validate", "fabro-workflows", @@ -1706,6 +1707,30 @@ dependencies = [ "tracing", ] +[[package]] +name = "fabro-telemetry" +version = "0.174.0" +dependencies = [ + "anyhow", + "base64", + "chrono", + "dirs", + "exec", + "fork", + "git2", + "insta", + "mac_address", + "md5", + "regex", + "reqwest 0.12.28", + "sentry", + "serde", + "serde_json", + "tokio", + "tracing", + "uuid", +] + [[package]] name = "fabro-tracker" version = "0.174.0" @@ -1734,19 +1759,10 @@ version = "0.174.0" dependencies = [ "aho-corasick", "anyhow", - "base64", - "chrono", "console 0.15.11", "dirs", - "exec", - "fork", - "git2", "insta", - "mac_address", - "md5", "regex", - "reqwest 0.12.28", - "sentry", "serde", "serde_json", "tempfile", @@ -1755,7 +1771,6 @@ dependencies = [ "toml", "tracing", "tracing-subscriber", - "uuid", ] [[package]] diff --git a/lib/crates/fabro-cli/Cargo.toml b/lib/crates/fabro-cli/Cargo.toml index f00ec69e8..adc4eb6bf 100644 --- a/lib/crates/fabro-cli/Cargo.toml +++ b/lib/crates/fabro-cli/Cargo.toml @@ -35,6 +35,7 @@ fabro-validate = { path = "../fabro-validate" } fabro-workflows = { path = "../fabro-workflows" } fabro-api = { path = "../fabro-api", optional = true } fabro-beastie = { path = "../fabro-beastie", optional = true } +fabro-telemetry = { path = "../fabro-telemetry" } fabro-util = { path = "../fabro-util" } clap.workspace = true cli-table.workspace = true diff --git a/lib/crates/fabro-cli/src/main.rs b/lib/crates/fabro-cli/src/main.rs index 8eef1ab42..e0b24994a 100644 --- a/lib/crates/fabro-cli/src/main.rs +++ b/lib/crates/fabro-cli/src/main.rs @@ -376,7 +376,8 @@ fn detach_run(args: commands::run::RunArgs) -> Result<()> { #[tokio::main] async fn main() { - fabro_util::telemetry::panic::install_panic_hook(); + fabro_telemetry::panic::install_panic_hook(); + fabro_telemetry::init_cli(); let start = std::time::Instant::now(); let raw_args: Vec = std::env::args().collect(); @@ -384,7 +385,29 @@ async fn main() { let (command_name, result) = main_inner().await; let duration_ms = start.elapsed().as_millis() as u64; - send_telemetry_event(&raw_args, &command_name, duration_ms, &result); + let is_error = result.is_err(); + if is_error { + fabro_telemetry::track!("CLI Errored", { + "subcommand": command_name, + "command": fabro_telemetry::sanitize::sanitize_command(&raw_args, &command_name), + "durationMs": duration_ms, + "repository": fabro_telemetry::git::repository_identifier(), + "ci": std::env::var("CI").is_ok(), + "success": false, + "exitCode": 1, + }, error); + } else { + fabro_telemetry::track!("CLI Executed", { + "subcommand": command_name, + "command": fabro_telemetry::sanitize::sanitize_command(&raw_args, &command_name), + "durationMs": duration_ms, + "repository": fabro_telemetry::git::repository_identifier(), + "ci": std::env::var("CI").is_ok(), + "success": true, + "exitCode": 0, + }); + } + fabro_telemetry::shutdown(); if let Err(err) = result { let style = console::Style::new().red().bold(); @@ -408,48 +431,6 @@ async fn main() { } } -fn send_telemetry_event( - raw_args: &[String], - command_name: &str, - duration_ms: u64, - result: &Result<()>, -) { - let is_error = result.is_err(); - let telemetry = match fabro_util::telemetry::Telemetry::for_cli() { - Ok(t) => t, - Err(err) => { - debug!(%err, "Telemetry initialization failed"); - return; - } - }; - if !telemetry.should_track(is_error) { - return; - } - - let event_name = if is_error { - "Command Error" - } else { - "Command Run" - }; - let properties = serde_json::json!({ - "subcommand": command_name, - "command": fabro_util::telemetry::sanitize::sanitize_command(raw_args, command_name), - "durationMs": duration_ms, - "repository": fabro_util::telemetry::git::repository_identifier(), - "ci": std::env::var("CI").is_ok(), - "success": !is_error, - "exitCode": if is_error { 1 } else { 0 }, - }); - - let track = telemetry.build_track(event_name, properties); - fabro_util::telemetry::sender::emit(&[track]); - debug!( - event = event_name, - subcommand = command_name, - "Telemetry event queued" - ); -} - async fn main_inner() -> (String, Result<()>) { let _ = rustls::crypto::ring::default_provider().install_default(); @@ -932,12 +913,12 @@ async fn main_inner() -> (String, Result<()>) { } }, Command::SendAnalytics { path } => { - let result = fabro_util::telemetry::sender::upload(&path).await; + let result = fabro_telemetry::sender::upload(&path).await; let _ = std::fs::remove_file(&path); result?; } Command::SendPanic { path } => { - let result = fabro_util::telemetry::panic::send_panic_to_sentry(&path).await; + let result = fabro_telemetry::panic::capture(&path).await; let _ = std::fs::remove_file(&path); result?; } diff --git a/lib/crates/fabro-telemetry/Cargo.toml b/lib/crates/fabro-telemetry/Cargo.toml new file mode 100644 index 000000000..201ee48f9 --- /dev/null +++ b/lib/crates/fabro-telemetry/Cargo.toml @@ -0,0 +1,31 @@ +[package] +name = "fabro-telemetry" +edition.workspace = true +version.workspace = true +license.workspace = true +description = "Telemetry, analytics, and crash reporting" + +[lib] +doctest = false + +[dependencies] +anyhow.workspace = true +base64.workspace = true +chrono.workspace = true +dirs.workspace = true +exec.workspace = true +fork.workspace = true +git2.workspace = true +mac_address.workspace = true +md5.workspace = true +regex.workspace = true +reqwest = { workspace = true, features = ["blocking"] } +sentry.workspace = true +serde.workspace = true +serde_json.workspace = true +tokio.workspace = true +tracing.workspace = true +uuid.workspace = true + +[dev-dependencies] +insta = { workspace = true } diff --git a/lib/crates/fabro-util/src/telemetry/anonymous_id.rs b/lib/crates/fabro-telemetry/src/anonymous_id.rs similarity index 100% rename from lib/crates/fabro-util/src/telemetry/anonymous_id.rs rename to lib/crates/fabro-telemetry/src/anonymous_id.rs diff --git a/lib/crates/fabro-telemetry/src/buffer.rs b/lib/crates/fabro-telemetry/src/buffer.rs new file mode 100644 index 000000000..caf717afb --- /dev/null +++ b/lib/crates/fabro-telemetry/src/buffer.rs @@ -0,0 +1,221 @@ +use std::sync::mpsc::Receiver; +use std::time::{Duration, Instant}; + +use crate::event::Track; + +pub(crate) struct BufferPolicy { + pub count_threshold: usize, + pub time_threshold: Duration, +} + +impl Default for BufferPolicy { + fn default() -> Self { + Self { + count_threshold: 20, + time_threshold: Duration::from_secs(60), + } + } +} + +pub(crate) fn consumer_loop( + rx: Receiver, + config: BufferPolicy, + mid_flush: impl Fn(&[Track]), + final_flush: impl Fn(&[Track]), +) { + let mut buffer: Vec = Vec::new(); + let mut next_flush = Instant::now() + config.time_threshold; + + loop { + let timeout = next_flush.saturating_duration_since(Instant::now()); + match rx.recv_timeout(timeout) { + Ok(track) => { + buffer.push(track); + if buffer.len() >= config.count_threshold { + mid_flush(&buffer); + buffer.clear(); + next_flush = Instant::now() + config.time_threshold; + } + } + Err(std::sync::mpsc::RecvTimeoutError::Timeout) => { + if !buffer.is_empty() { + mid_flush(&buffer); + buffer.clear(); + } + next_flush = Instant::now() + config.time_threshold; + } + Err(std::sync::mpsc::RecvTimeoutError::Disconnected) => { + // Drain any remaining events + while let Ok(track) = rx.try_recv() { + buffer.push(track); + } + if !buffer.is_empty() { + final_flush(&buffer); + } + return; + } + } + } +} + +#[cfg(test)] +mod tests { + use super::*; + use crate::event::User; + use serde_json::json; + use std::sync::{mpsc, Arc, Mutex}; + + fn make_track(event: &str) -> Track { + Track { + user: User::AnonymousId { + anonymous_id: "test".to_string(), + }, + event: event.to_string(), + properties: json!({}), + context: None, + timestamp: None, + message_id: format!("msg-{event}"), + } + } + + #[test] + fn flushes_on_count_threshold() { + let (tx, rx) = mpsc::channel(); + let mid_flushes: Arc>>> = Arc::new(Mutex::new(Vec::new())); + let final_flushes: Arc>>> = Arc::new(Mutex::new(Vec::new())); + + let mid = mid_flushes.clone(); + let fin = final_flushes.clone(); + + tx.send(make_track("e1")).unwrap(); + tx.send(make_track("e2")).unwrap(); + drop(tx); + + consumer_loop( + rx, + BufferPolicy { + count_threshold: 2, + time_threshold: Duration::from_secs(60), + }, + move |tracks| { + let events: Vec = tracks.iter().map(|t| t.event.clone()).collect(); + mid.lock().unwrap().push(events); + }, + move |tracks| { + let events: Vec = tracks.iter().map(|t| t.event.clone()).collect(); + fin.lock().unwrap().push(events); + }, + ); + + let mid = mid_flushes.lock().unwrap(); + assert_eq!(mid.len(), 1); + assert_eq!(mid[0], vec!["e1", "e2"]); + + let fin = final_flushes.lock().unwrap(); + assert!(fin.is_empty()); + } + + #[test] + fn no_flush_when_empty() { + let (tx, rx) = mpsc::channel::(); + let mid_called = Arc::new(Mutex::new(false)); + let final_called = Arc::new(Mutex::new(false)); + + let mid = mid_called.clone(); + let fin = final_called.clone(); + + drop(tx); + + consumer_loop( + rx, + BufferPolicy { + count_threshold: 2, + time_threshold: Duration::from_secs(60), + }, + move |_| { + *mid.lock().unwrap() = true; + }, + move |_| { + *fin.lock().unwrap() = true; + }, + ); + + assert!(!*mid_called.lock().unwrap()); + assert!(!*final_called.lock().unwrap()); + } + + #[test] + fn flushes_on_time_threshold() { + let (tx, rx) = mpsc::channel(); + let mid_flushes: Arc>>> = Arc::new(Mutex::new(Vec::new())); + let final_flushes: Arc>>> = Arc::new(Mutex::new(Vec::new())); + + let mid = mid_flushes.clone(); + let fin = final_flushes.clone(); + + tx.send(make_track("e1")).unwrap(); + // Don't drop yet — let time threshold fire + let handle = std::thread::spawn(move || { + consumer_loop( + rx, + BufferPolicy { + count_threshold: 100, // won't trigger + time_threshold: Duration::from_millis(50), + }, + move |tracks| { + let events: Vec = tracks.iter().map(|t| t.event.clone()).collect(); + mid.lock().unwrap().push(events); + }, + move |tracks| { + let events: Vec = tracks.iter().map(|t| t.event.clone()).collect(); + fin.lock().unwrap().push(events); + }, + ); + }); + + // Wait for time threshold to fire, then drop sender + std::thread::sleep(Duration::from_millis(150)); + drop(tx); + handle.join().unwrap(); + + let mid = mid_flushes.lock().unwrap(); + assert!(!mid.is_empty(), "time-based flush should have fired"); + assert_eq!(mid[0], vec!["e1"]); + } + + #[test] + fn flushes_remaining_on_disconnect() { + let (tx, rx) = mpsc::channel(); + let mid_flushes: Arc>>> = Arc::new(Mutex::new(Vec::new())); + let final_flushes: Arc>>> = Arc::new(Mutex::new(Vec::new())); + + let mid = mid_flushes.clone(); + let fin = final_flushes.clone(); + + tx.send(make_track("e1")).unwrap(); + drop(tx); // disconnect immediately, below count threshold + + consumer_loop( + rx, + BufferPolicy { + count_threshold: 100, // won't trigger + time_threshold: Duration::from_secs(60), + }, + move |tracks| { + let events: Vec = tracks.iter().map(|t| t.event.clone()).collect(); + mid.lock().unwrap().push(events); + }, + move |tracks| { + let events: Vec = tracks.iter().map(|t| t.event.clone()).collect(); + fin.lock().unwrap().push(events); + }, + ); + + let mid = mid_flushes.lock().unwrap(); + assert!(mid.is_empty()); + + let fin = final_flushes.lock().unwrap(); + assert_eq!(fin.len(), 1); + assert_eq!(fin[0], vec!["e1"]); + } +} diff --git a/lib/crates/fabro-util/src/telemetry/context.rs b/lib/crates/fabro-telemetry/src/context.rs similarity index 100% rename from lib/crates/fabro-util/src/telemetry/context.rs rename to lib/crates/fabro-telemetry/src/context.rs diff --git a/lib/crates/fabro-util/src/telemetry/event.rs b/lib/crates/fabro-telemetry/src/event.rs similarity index 100% rename from lib/crates/fabro-util/src/telemetry/event.rs rename to lib/crates/fabro-telemetry/src/event.rs diff --git a/lib/crates/fabro-util/src/telemetry/git.rs b/lib/crates/fabro-telemetry/src/git.rs similarity index 100% rename from lib/crates/fabro-util/src/telemetry/git.rs rename to lib/crates/fabro-telemetry/src/git.rs diff --git a/lib/crates/fabro-telemetry/src/lib.rs b/lib/crates/fabro-telemetry/src/lib.rs new file mode 100644 index 000000000..17ef30d6b --- /dev/null +++ b/lib/crates/fabro-telemetry/src/lib.rs @@ -0,0 +1,245 @@ +pub mod anonymous_id; +mod buffer; +pub mod context; +pub mod event; +pub mod git; +pub mod panic; +pub mod sanitize; +pub mod sender; +pub mod spawn; + +use std::sync::mpsc; +use std::sync::{Mutex, OnceLock}; +use std::thread::JoinHandle; + +use chrono::Utc; +use serde_json::Value; +use uuid::Uuid; + +use event::{Track, User}; + +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum TelemetryLevel { + Off, + Errors, + All, +} + +struct Global { + sender: Mutex>>, + anonymous_id: String, + context: Value, + level: TelemetryLevel, + thread: Mutex>>, +} + +static GLOBAL: OnceLock = OnceLock::new(); + +/// Initialize telemetry for the CLI. Spawns a background thread for buffered delivery. +/// No-op if telemetry level is `Off`. +pub fn init_cli() { + let level = telemetry_level(); + if level == TelemetryLevel::Off { + return; + } + + let anonymous_id = match anonymous_id::compute_cli_id() { + Ok(id) => id, + Err(err) => { + tracing::debug!(%err, "telemetry: failed to compute CLI anonymous id"); + return; + } + }; + + init_inner(level, anonymous_id); +} + +/// Initialize telemetry for the server. Spawns a background thread for buffered delivery. +/// No-op if telemetry level is `Off`. +pub fn init_server() { + let level = telemetry_level(); + if level == TelemetryLevel::Off { + return; + } + + let anonymous_id = match anonymous_id::load_or_create_server_id() { + Ok(id) => id, + Err(err) => { + tracing::debug!(%err, "telemetry: failed to load/create server anonymous id"); + return; + } + }; + + init_inner(level, anonymous_id); +} + +fn init_inner(level: TelemetryLevel, anonymous_id: String) { + let ctx = context::build_context(); + let (tx, rx) = mpsc::channel(); + + let handle = std::thread::Builder::new() + .name("telemetry".to_string()) + .spawn(move || { + buffer::consumer_loop( + rx, + buffer::BufferPolicy::default(), + |tracks| { + if let Err(err) = sender::upload_blocking(tracks) { + tracing::debug!(%err, "telemetry: mid-run flush failed"); + } + }, + |tracks| { + sender::emit(tracks); + }, + ); + }) + .expect("failed to spawn telemetry thread"); + + let _ = GLOBAL.set(Global { + sender: Mutex::new(Some(tx)), + anonymous_id, + context: ctx, + level, + thread: Mutex::new(Some(handle)), + }); +} + +/// Shut down telemetry: close the channel and wait for the background thread to finish. +/// The final flush uses the detached-subprocess pattern so events survive process exit. +pub fn shutdown() { + let Some(global) = GLOBAL.get() else { + return; + }; + + // Drop sender to close the channel + if let Ok(mut sender) = global.sender.lock() { + sender.take(); + } + + // Join the thread (blocks until final flush subprocess is spawned) + if let Ok(mut handle) = global.thread.lock() { + if let Some(h) = handle.take() { + let _ = h.join(); + } + } +} + +/// Check whether an event should be tracked given a telemetry level. +pub fn should_track_for_level(level: TelemetryLevel, is_error: bool) -> bool { + match level { + TelemetryLevel::Off => false, + TelemetryLevel::Errors => is_error, + TelemetryLevel::All => true, + } +} + +/// Internal function called by the `track!` macro. Do not call directly. +#[doc(hidden)] +pub fn _track_inner(event: &str, properties: Value, is_error: bool) { + let Some(global) = GLOBAL.get() else { + return; + }; + + if !should_track_for_level(global.level, is_error) { + return; + } + + let track = Track { + user: User::AnonymousId { + anonymous_id: global.anonymous_id.clone(), + }, + event: event.to_string(), + properties, + context: Some(global.context.clone()), + timestamp: Some(Utc::now().to_rfc3339()), + message_id: Uuid::new_v4().to_string(), + }; + + if let Ok(sender) = global.sender.lock() { + if let Some(tx) = sender.as_ref() { + let _ = tx.send(track); + } + } +} + +/// Record a telemetry event. Non-blocking — queues to the background buffer. +/// +/// # Examples +/// +/// ```ignore +/// fabro_telemetry::track!("Command Run", { +/// "subcommand": "run", +/// "durationMs": 1234, +/// }); +/// +/// fabro_telemetry::track!("Command Error", { "subcommand": "run" }, error); +/// ``` +#[macro_export] +macro_rules! track { + ($event:expr, { $($tt:tt)* }) => { + $crate::_track_inner($event, ::serde_json::json!({ $($tt)* }), false) + }; + ($event:expr, { $($tt:tt)* }, error) => { + $crate::_track_inner($event, ::serde_json::json!({ $($tt)* }), true) + }; +} + +pub fn telemetry_level() -> TelemetryLevel { + telemetry_level_from(std::env::var("FABRO_TELEMETRY").ok().as_deref()) +} + +pub fn telemetry_level_from(env_value: Option<&str>) -> TelemetryLevel { + match env_value { + Some("off") => TelemetryLevel::Off, + Some("errors") => TelemetryLevel::Errors, + Some("all") => TelemetryLevel::All, + _ => { + if cfg!(debug_assertions) { + TelemetryLevel::Off + } else { + TelemetryLevel::All + } + } + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn telemetry_level_defaults_to_off_in_debug() { + assert_eq!(telemetry_level_from(None), TelemetryLevel::Off); + } + + #[test] + fn telemetry_level_parses_env_var() { + assert_eq!(telemetry_level_from(Some("all")), TelemetryLevel::All); + assert_eq!(telemetry_level_from(Some("errors")), TelemetryLevel::Errors); + assert_eq!(telemetry_level_from(Some("off")), TelemetryLevel::Off); + } + + #[test] + fn should_track_for_level_off() { + assert!(!should_track_for_level(TelemetryLevel::Off, false)); + assert!(!should_track_for_level(TelemetryLevel::Off, true)); + } + + #[test] + fn should_track_for_level_errors() { + assert!(!should_track_for_level(TelemetryLevel::Errors, false)); + assert!(should_track_for_level(TelemetryLevel::Errors, true)); + } + + #[test] + fn should_track_for_level_all() { + assert!(should_track_for_level(TelemetryLevel::All, false)); + assert!(should_track_for_level(TelemetryLevel::All, true)); + } + + #[test] + fn track_inner_noop_when_not_initialized() { + // GLOBAL is not set in unit tests, so this should silently return + _track_inner("Test Event", serde_json::json!({"key": "value"}), false); + } +} diff --git a/lib/crates/fabro-util/src/telemetry/panic.rs b/lib/crates/fabro-telemetry/src/panic.rs similarity index 85% rename from lib/crates/fabro-util/src/telemetry/panic.rs rename to lib/crates/fabro-telemetry/src/panic.rs index ff52a593e..065c51f44 100644 --- a/lib/crates/fabro-util/src/telemetry/panic.rs +++ b/lib/crates/fabro-telemetry/src/panic.rs @@ -3,7 +3,7 @@ use std::path::Path; use sentry::protocol::{Event, Exception, Mechanism}; -use super::TelemetryLevel; +use crate::TelemetryLevel; const SENTRY_DSN: Option<&str> = option_env!("SENTRY_DSN"); @@ -20,7 +20,7 @@ pub fn install_panic_hook() { } /// Build a Sentry event from panic info. Exposed for testing. -pub fn build_panic_event(message: &str) -> Event<'static> { +pub fn build_event(message: &str) -> Event<'static> { let mut event = Event::new(); event.level = sentry::Level::Fatal; @@ -52,7 +52,7 @@ pub fn build_panic_event(message: &str) -> Event<'static> { ); // Set release to the package version. - event.release = Some(crate::version::FABRO_VERSION.into()); + event.release = Some(env!("CARGO_PKG_VERSION").into()); event } @@ -80,7 +80,7 @@ fn report_panic(info: &PanicHookInfo<'_>) { return; } - let level = super::telemetry_level(); + let level = crate::telemetry_level(); if level == TelemetryLevel::Off { return; } @@ -90,7 +90,7 @@ fn report_panic(info: &PanicHookInfo<'_>) { return; } - let event = build_panic_event(&message); + let event = build_event(&message); spawn_panic_sender(event); } @@ -102,14 +102,14 @@ fn spawn_panic_sender(event: Event<'static>) { }; let filename = format!("fabro-panic-{}.json", event.event_id); - super::spawn::spawn_fabro_subcommand("__send_panic", &filename, &json); + crate::spawn::spawn_fabro_subcommand("__send_panic", &filename, &json); } /// Send a serialized Sentry panic event. Called by the `__send_panic` subcommand. /// /// Reads the JSON event from `path` and sends it to Sentry. /// No-ops if `SENTRY_DSN` was not set at compile time. -pub async fn send_panic_to_sentry(path: &Path) -> anyhow::Result<()> { +pub async fn capture(path: &Path) -> anyhow::Result<()> { let dsn = SENTRY_DSN.ok_or_else(|| anyhow::anyhow!("SENTRY_DSN not set at compile time"))?; let json = std::fs::read(path)?; @@ -130,8 +130,8 @@ mod tests { use super::*; #[test] - fn build_panic_event_structure() { - let event = build_panic_event("test panic message"); + fn build_event_structure() { + let event = build_event("test panic message"); assert_eq!(event.level, sentry::Level::Fatal); assert_eq!(event.exception.values.len(), 1); @@ -163,12 +163,9 @@ mod tests { #[test] fn report_panic_noop_when_telemetry_off() { - use crate::env::TestEnv; - use std::collections::HashMap; - // Verify telemetry_level_from returns Off without mutating process env. - let env = TestEnv(HashMap::from([("FABRO_TELEMETRY".into(), "off".into())])); + // Verify telemetry_level_from returns Off for "off". assert_eq!( - super::super::telemetry_level_from(&env), + crate::telemetry_level_from(Some("off")), TelemetryLevel::Off ); } @@ -177,7 +174,7 @@ mod tests { fn send_panic_noops_without_dsn() { // SENTRY_DSN is not set at compile time in tests, so this should error. let rt = tokio::runtime::Runtime::new().unwrap(); - let result = rt.block_on(send_panic_to_sentry(Path::new("/nonexistent"))); + let result = rt.block_on(capture(Path::new("/nonexistent"))); assert!(result.is_err()); let err_msg = result.unwrap_err().to_string(); assert!(err_msg.contains("SENTRY_DSN not set")); @@ -185,7 +182,7 @@ mod tests { #[test] fn event_round_trips_through_json() { - let event = build_panic_event("roundtrip test"); + let event = build_event("roundtrip test"); let json = serde_json::to_vec(&event).unwrap(); let deserialized: Event<'static> = serde_json::from_slice(&json).unwrap(); assert_eq!(deserialized.level, sentry::Level::Fatal); diff --git a/lib/crates/fabro-util/src/telemetry/sanitize.rs b/lib/crates/fabro-telemetry/src/sanitize.rs similarity index 100% rename from lib/crates/fabro-util/src/telemetry/sanitize.rs rename to lib/crates/fabro-telemetry/src/sanitize.rs diff --git a/lib/crates/fabro-telemetry/src/sender.rs b/lib/crates/fabro-telemetry/src/sender.rs index 917eaf726..dd39ca14d 100644 --- a/lib/crates/fabro-telemetry/src/sender.rs +++ b/lib/crates/fabro-telemetry/src/sender.rs @@ -75,6 +75,43 @@ fn build_segment_batch(content: &str) -> Option { Some(serde_json::json!({ "batch": batch })) } +/// Sends track events to Segment synchronously. Used for mid-run flushes +/// on the background telemetry thread. +/// No-ops if `SEGMENT_WRITE_KEY` was not set at compile time or `tracks` is empty. +pub fn upload_blocking(tracks: &[Track]) -> anyhow::Result<()> { + let write_key = SEGMENT_WRITE_KEY + .ok_or_else(|| anyhow::anyhow!("SEGMENT_WRITE_KEY not set at compile time"))?; + + if tracks.is_empty() { + return Ok(()); + } + + let lines: Vec = tracks + .iter() + .filter_map(|t| serde_json::to_string(t).ok()) + .collect(); + + let content = lines.join("\n"); + let payload = match build_segment_batch(&content) { + Some(p) => p, + None => return Ok(()), + }; + + let auth = STANDARD.encode(format!("{write_key}:")); + + let resp = reqwest::blocking::Client::new() + .post(format!("{SEGMENT_BASE_URL}/v1/batch")) + .header("Authorization", format!("Basic {auth}")) + .json(&payload) + .send()?; + + if !resp.status().is_success() { + anyhow::bail!("segment API returned status {}", resp.status()); + } + + Ok(()) +} + /// Reads a JSONL file of serialized track events from `path` and sends them /// to Segment as a batch. /// Called by the `__send_analytics` subcommand. @@ -208,7 +245,35 @@ mod tests { emit(&[]); } - // -- Step 3: upload() tests -- + // -- Step 3: upload_blocking() tests -- + + #[test] + fn upload_blocking_noops_without_write_key() { + let track = Track { + user: User::AnonymousId { + anonymous_id: "test".to_string(), + }, + event: "test".to_string(), + properties: json!({}), + context: None, + timestamp: None, + message_id: "msg-test".to_string(), + }; + + let result = upload_blocking(&[track]); + assert!(result.is_err()); + let err_msg = result.unwrap_err().to_string(); + assert!(err_msg.contains("SEGMENT_WRITE_KEY not set")); + } + + #[test] + fn upload_blocking_noops_with_empty_tracks() { + // With no write key, but empty tracks should still error on write key check + let result = upload_blocking(&[]); + assert!(result.is_err()); + } + + // -- Step 4: upload() tests -- #[test] fn upload_noops_without_write_key() { diff --git a/lib/crates/fabro-util/src/telemetry/spawn.rs b/lib/crates/fabro-telemetry/src/spawn.rs similarity index 100% rename from lib/crates/fabro-util/src/telemetry/spawn.rs rename to lib/crates/fabro-telemetry/src/spawn.rs diff --git a/lib/crates/fabro-util/Cargo.toml b/lib/crates/fabro-util/Cargo.toml index 55c173d47..e7f416ba9 100644 --- a/lib/crates/fabro-util/Cargo.toml +++ b/lib/crates/fabro-util/Cargo.toml @@ -15,21 +15,11 @@ termimad.workspace = true aho-corasick.workspace = true serde_json.workspace = true serde.workspace = true -reqwest.workspace = true tokio.workspace = true -uuid.workspace = true -chrono.workspace = true -base64.workspace = true dirs.workspace = true tracing.workspace = true tracing-subscriber.workspace = true anyhow.workspace = true -mac_address.workspace = true -md5.workspace = true -git2.workspace = true -sentry.workspace = true -fork.workspace = true -exec.workspace = true [build-dependencies] toml = "0.8" diff --git a/lib/crates/fabro-util/src/lib.rs b/lib/crates/fabro-util/src/lib.rs index 5a6880219..1dc496e23 100644 --- a/lib/crates/fabro-util/src/lib.rs +++ b/lib/crates/fabro-util/src/lib.rs @@ -3,6 +3,5 @@ pub mod env; pub mod path; pub mod redact; pub mod run_log; -pub mod telemetry; pub mod terminal; pub mod version; diff --git a/lib/crates/fabro-util/src/telemetry/mod.rs b/lib/crates/fabro-util/src/telemetry/mod.rs deleted file mode 100644 index 3d4e7f88a..000000000 --- a/lib/crates/fabro-util/src/telemetry/mod.rs +++ /dev/null @@ -1,165 +0,0 @@ -pub mod anonymous_id; -pub mod context; -pub mod event; -pub mod git; -pub mod panic; -pub mod sanitize; -pub mod sender; -pub mod spawn; - -use anyhow::Result; -use chrono::Utc; -use serde_json::Value; -use uuid::Uuid; - -use event::{Track, User}; - -#[derive(Debug, Clone, Copy, PartialEq, Eq)] -pub enum TelemetryLevel { - Off, - Errors, - All, -} - -pub struct Telemetry { - level: TelemetryLevel, - anonymous_id: String, - context: Value, -} - -impl Telemetry { - fn new(anonymous_id: String) -> Self { - Self { - level: telemetry_level(), - anonymous_id, - context: context::build_context(), - } - } - - pub fn for_server() -> Result { - Ok(Self::new(anonymous_id::load_or_create_server_id()?)) - } - - pub fn for_cli() -> Result { - Ok(Self::new(anonymous_id::compute_cli_id()?)) - } - - pub fn level(&self) -> &TelemetryLevel { - &self.level - } - - pub fn should_track(&self, is_error: bool) -> bool { - match self.level { - TelemetryLevel::Off => false, - TelemetryLevel::Errors => is_error, - TelemetryLevel::All => true, - } - } - - pub fn build_track(&self, event: &str, properties: Value) -> Track { - Track { - user: User::AnonymousId { - anonymous_id: self.anonymous_id.clone(), - }, - event: event.to_string(), - properties, - context: Some(self.context.clone()), - timestamp: Some(Utc::now().to_rfc3339()), - message_id: Uuid::new_v4().to_string(), - } - } -} - -pub fn telemetry_level() -> TelemetryLevel { - telemetry_level_from(&crate::env::SystemEnv) -} - -pub fn telemetry_level_from(env: &dyn crate::env::Env) -> TelemetryLevel { - match env.var("FABRO_TELEMETRY").as_deref() { - Ok("off") => TelemetryLevel::Off, - Ok("errors") => TelemetryLevel::Errors, - Ok("all") => TelemetryLevel::All, - _ => { - if cfg!(debug_assertions) { - TelemetryLevel::Off - } else { - TelemetryLevel::All - } - } - } -} - -#[cfg(test)] -mod tests { - use super::*; - use crate::env::TestEnv; - use serde_json::json; - use std::collections::HashMap; - - #[test] - fn telemetry_level_defaults_to_off_in_debug() { - // In test builds (debug_assertions=true), default is Off - let env = TestEnv(HashMap::new()); - assert_eq!(telemetry_level_from(&env), TelemetryLevel::Off); - } - - #[test] - fn telemetry_level_parses_env_var() { - let env = TestEnv(HashMap::from([("FABRO_TELEMETRY".into(), "all".into())])); - assert_eq!(telemetry_level_from(&env), TelemetryLevel::All); - - let env = TestEnv(HashMap::from([("FABRO_TELEMETRY".into(), "errors".into())])); - assert_eq!(telemetry_level_from(&env), TelemetryLevel::Errors); - - let env = TestEnv(HashMap::from([("FABRO_TELEMETRY".into(), "off".into())])); - assert_eq!(telemetry_level_from(&env), TelemetryLevel::Off); - } - - #[test] - fn should_track_respects_level() { - let telemetry = Telemetry { - level: TelemetryLevel::Off, - anonymous_id: "test".to_string(), - context: json!({}), - }; - assert!(!telemetry.should_track(false)); - assert!(!telemetry.should_track(true)); - - let telemetry = Telemetry { - level: TelemetryLevel::Errors, - anonymous_id: "test".to_string(), - context: json!({}), - }; - assert!(!telemetry.should_track(false)); - assert!(telemetry.should_track(true)); - - let telemetry = Telemetry { - level: TelemetryLevel::All, - anonymous_id: "test".to_string(), - context: json!({}), - }; - assert!(telemetry.should_track(false)); - assert!(telemetry.should_track(true)); - } - - #[test] - fn build_track_populates_fields() { - let telemetry = Telemetry { - level: TelemetryLevel::All, - anonymous_id: "anon-123".to_string(), - context: json!({"app": {"name": "fabro"}}), - }; - - let track = telemetry.build_track("Test Event", json!({"key": "value"})); - assert_eq!(track.event, "Test Event"); - assert_eq!(track.properties["key"], "value"); - assert!(track.context.is_some()); - assert!(track.timestamp.is_some()); - assert!(!track.message_id.is_empty()); - - match &track.user { - User::AnonymousId { anonymous_id } => assert_eq!(anonymous_id, "anon-123"), - _ => panic!("expected AnonymousId variant"), - } - } -} diff --git a/lib/crates/fabro-util/src/telemetry/sender.rs b/lib/crates/fabro-util/src/telemetry/sender.rs deleted file mode 100644 index 48025cab4..000000000 --- a/lib/crates/fabro-util/src/telemetry/sender.rs +++ /dev/null @@ -1,219 +0,0 @@ -use std::path::Path; - -use base64::engine::general_purpose::STANDARD; -use base64::Engine; -use uuid::Uuid; - -use super::event::Track; - -const SEGMENT_API_URL: &str = "https://api.segment.io/v1/batch"; -const SEGMENT_WRITE_KEY: Option<&str> = option_env!("SEGMENT_WRITE_KEY"); - -/// Serializes the track events as JSONL to a temp file and spawns a detached -/// subprocess (`fabro __send_analytics `) to deliver them. This ensures -/// the events are sent even if the parent CLI process exits immediately. -/// -/// No-ops if the SEGMENT_WRITE_KEY was not set at compile time or `tracks` is empty. -pub fn emit(tracks: &[Track]) { - if SEGMENT_WRITE_KEY.is_none() { - tracing::debug!("telemetry: no SEGMENT_WRITE_KEY, skipping emit"); - return; - } - - if tracks.is_empty() { - return; - } - - spawn_sender(tracks); -} - -fn spawn_sender(tracks: &[Track]) { - let lines: Vec = tracks - .iter() - .filter_map(|t| serde_json::to_string(t).ok()) - .collect(); - - if lines.is_empty() { - return; - } - - let jsonl = lines.join("\n"); - let filename = format!("fabro-events-{}.jsonl", Uuid::new_v4()); - super::spawn::spawn_fabro_subcommand("__send_analytics", &filename, jsonl.as_bytes()); -} - -/// Parse JSONL content into a Segment batch payload. -/// -/// Each non-empty line is parsed as JSON, has `"type": "track"` injected, -/// and is collected into a `{"batch": [...]}` wrapper. -/// Returns `None` if no valid events are found. -fn build_segment_batch(content: &str) -> Option { - let mut batch = Vec::new(); - for line in content.lines() { - let line = line.trim(); - if line.is_empty() { - continue; - } - match serde_json::from_str::>(line) { - Ok(mut map) => { - map.insert("type".into(), "track".into()); - batch.push(serde_json::Value::Object(map)); - } - Err(err) => { - tracing::warn!(%err, "skipping malformed JSONL line"); - } - } - } - - if batch.is_empty() { - return None; - } - - Some(serde_json::json!({ "batch": batch })) -} - -/// Reads a JSONL file of serialized track events from `path` and sends them -/// to Segment as a batch. -/// Called by the `__send_analytics` subcommand. -/// No-ops if `SEGMENT_WRITE_KEY` was not set at compile time. -pub async fn upload(path: &Path) -> anyhow::Result<()> { - let write_key = SEGMENT_WRITE_KEY - .ok_or_else(|| anyhow::anyhow!("SEGMENT_WRITE_KEY not set at compile time"))?; - - let content = std::fs::read_to_string(path)?; - let payload = match build_segment_batch(&content) { - Some(p) => p, - None => return Ok(()), - }; - - let auth = STANDARD.encode(format!("{write_key}:")); - - let resp = reqwest::Client::new() - .post(SEGMENT_API_URL) - .header("Authorization", format!("Basic {auth}")) - .json(&payload) - .send() - .await?; - - if !resp.status().is_success() { - anyhow::bail!("segment API returned status {}", resp.status()); - } - - Ok(()) -} - -#[cfg(test)] -mod tests { - use super::*; - use crate::telemetry::event::User; - use serde_json::json; - - // -- Step 1: build_segment_batch tests -- - - #[test] - fn build_segment_batch_empty_content() { - assert!(build_segment_batch("").is_none()); - } - - #[test] - fn build_segment_batch_single_event() { - let line = r#"{"anonymousId":"abc","event":"Test","properties":{},"messageId":"m1"}"#; - let result = build_segment_batch(line).unwrap(); - - let batch = result["batch"].as_array().unwrap(); - assert_eq!(batch.len(), 1); - assert_eq!(batch[0]["type"], "track"); - assert_eq!(batch[0]["event"], "Test"); - assert_eq!(batch[0]["anonymousId"], "abc"); - } - - #[test] - fn build_segment_batch_multiple_events() { - let content = concat!( - r#"{"anonymousId":"a","event":"E1","properties":{},"messageId":"m1"}"#, - "\n", - r#"{"anonymousId":"b","event":"E2","properties":{},"messageId":"m2"}"#, - ); - let result = build_segment_batch(content).unwrap(); - - let batch = result["batch"].as_array().unwrap(); - assert_eq!(batch.len(), 2); - assert_eq!(batch[0]["type"], "track"); - assert_eq!(batch[0]["event"], "E1"); - assert_eq!(batch[1]["type"], "track"); - assert_eq!(batch[1]["event"], "E2"); - } - - #[test] - fn build_segment_batch_skips_malformed_lines() { - let content = concat!( - r#"{"anonymousId":"a","event":"Good","properties":{},"messageId":"m1"}"#, - "\n", - "this is not json", - ); - let result = build_segment_batch(content).unwrap(); - - let batch = result["batch"].as_array().unwrap(); - assert_eq!(batch.len(), 1); - assert_eq!(batch[0]["event"], "Good"); - } - - #[test] - fn build_segment_batch_all_malformed() { - let content = "not json\nalso not json\n"; - assert!(build_segment_batch(content).is_none()); - } - - #[test] - fn build_segment_batch_skips_blank_lines() { - let content = concat!( - "\n", - r#"{"anonymousId":"a","event":"E1","properties":{},"messageId":"m1"}"#, - "\n", - "\n", - ); - let result = build_segment_batch(content).unwrap(); - - let batch = result["batch"].as_array().unwrap(); - assert_eq!(batch.len(), 1); - } - - // -- Step 2: emit() tests -- - - #[test] - fn emit_noops_without_write_key() { - // SEGMENT_WRITE_KEY is not set at compile time in tests, - // so emit() should return immediately without spawning. - let track = Track { - user: User::AnonymousId { - anonymous_id: "test".to_string(), - }, - event: "test".to_string(), - properties: json!({}), - context: None, - timestamp: None, - message_id: "msg-test".to_string(), - }; - - // This should not panic or require a tokio runtime - // because it returns before reaching spawn - emit(&[track]); - } - - #[test] - fn emit_noops_with_empty_tracks() { - emit(&[]); - } - - // -- Step 3: upload() tests -- - - #[test] - fn upload_noops_without_write_key() { - // SEGMENT_WRITE_KEY is not set at compile time in tests, so this should error. - let rt = tokio::runtime::Runtime::new().unwrap(); - let result = rt.block_on(upload(Path::new("/nonexistent"))); - assert!(result.is_err()); - let err_msg = result.unwrap_err().to_string(); - assert!(err_msg.contains("SEGMENT_WRITE_KEY not set")); - } -}