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) <noreply@anthropic.com>
This commit is contained in:
Bryan Helmkamp 2026-03-19 14:15:33 -04:00
parent 1564303c2c
commit 94ecc0976b
No known key found for this signature in database
19 changed files with 631 additions and 469 deletions

View file

@ -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

35
Cargo.lock generated
View file

@ -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]]

View file

@ -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

View file

@ -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<String> = 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?;
}

View file

@ -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 }

View file

@ -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<Track>,
config: BufferPolicy,
mid_flush: impl Fn(&[Track]),
final_flush: impl Fn(&[Track]),
) {
let mut buffer: Vec<Track> = 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<Mutex<Vec<Vec<String>>>> = Arc::new(Mutex::new(Vec::new()));
let final_flushes: Arc<Mutex<Vec<Vec<String>>>> = 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<String> = tracks.iter().map(|t| t.event.clone()).collect();
mid.lock().unwrap().push(events);
},
move |tracks| {
let events: Vec<String> = 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::<Track>();
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<Mutex<Vec<Vec<String>>>> = Arc::new(Mutex::new(Vec::new()));
let final_flushes: Arc<Mutex<Vec<Vec<String>>>> = 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<String> = tracks.iter().map(|t| t.event.clone()).collect();
mid.lock().unwrap().push(events);
},
move |tracks| {
let events: Vec<String> = 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<Mutex<Vec<Vec<String>>>> = Arc::new(Mutex::new(Vec::new()));
let final_flushes: Arc<Mutex<Vec<Vec<String>>>> = 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<String> = tracks.iter().map(|t| t.event.clone()).collect();
mid.lock().unwrap().push(events);
},
move |tracks| {
let events: Vec<String> = 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"]);
}
}

View file

@ -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<Option<mpsc::Sender<Track>>>,
anonymous_id: String,
context: Value,
level: TelemetryLevel,
thread: Mutex<Option<JoinHandle<()>>>,
}
static GLOBAL: OnceLock<Global> = 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);
}
}

View file

@ -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);

View file

@ -75,6 +75,43 @@ fn build_segment_batch(content: &str) -> Option<serde_json::Value> {
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<String> = 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() {

View file

@ -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"

View file

@ -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;

View file

@ -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<Self> {
Ok(Self::new(anonymous_id::load_or_create_server_id()?))
}
pub fn for_cli() -> Result<Self> {
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"),
}
}
}

View file

@ -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 <path>`) 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<String> = 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<serde_json::Value> {
let mut batch = Vec::new();
for line in content.lines() {
let line = line.trim();
if line.is_empty() {
continue;
}
match serde_json::from_str::<serde_json::Map<String, serde_json::Value>>(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"));
}
}