#![deny(missing_docs)]
mod actor;
pub mod buffer;
pub mod client;
pub mod counters;
pub mod decision;
pub mod envelope;
pub mod event;
pub mod notice;
#[cfg(test)]
mod tests;
use std::sync::OnceLock;
use std::sync::atomic::{AtomicU8, Ordering};
use std::time::Duration;
pub use actor::{BATCH_MAX_BYTES, BATCH_MAX_EVENTS, FlushOutcome};
pub use counters::{Counter, ErrorCounter, SessionCounters};
pub use decision::{
EndpointError, TELEMETRY_DIR, TelemetryConsent, TelemetryDecision, decide, decide_in_home,
load_setup_state_for_decision, load_setup_state_for_decision_at, re_decide, validate_endpoint,
};
pub use envelope::reduce_panic_site;
pub use event::{
Arch, Batch, ColdStartBucket, Counters, DurationBucket, Errors, Event, ExitClass, InstallKind,
Libc, Os, SCHEMA_VERSION, SessionSource, Surface, TurnWall,
};
pub const SHUTDOWN_FLUSH_TIMEOUT: Duration = Duration::from_secs(3);
struct Armed {
handle: actor::Handle,
root: std::path::PathBuf,
exit_class: AtomicU8,
}
static ARMED: OnceLock<Armed> = OnceLock::new();
pub fn init(consent: TelemetryConsent) {
if ARMED.get().is_some() {
return;
}
let root = consent.root().to_path_buf();
let observed_generation = consent.tombstone_generation().cloned();
let config_path = consent.config_path().map(std::path::Path::to_path_buf);
if let Err(error) = buffer::arm(&root, observed_generation.as_ref(), || {
decision::permission_still_enabled(config_path.as_deref(), &root)
}) {
tracing::debug!("telemetry could not prepare its buffer: {error}");
return;
}
let context = actor::Context {
root: root.clone(),
endpoint: consent.endpoint().map(str::to_string),
surface: consent.surface(),
config_path: consent.config_path().map(std::path::Path::to_path_buf),
app_version: env!("CARGO_PKG_VERSION").to_string(),
git_sha: envelope::release_build_sha(),
tty: envelope::current_tty(),
};
let _ = ARMED.set(Armed {
handle: actor::Handle::spawn(context),
root: root.clone(),
exit_class: AtomicU8::new(ExitClass::Clean.as_u8()),
});
record_install_or_upgrade(&root);
}
fn record_install_or_upgrade(root: &std::path::Path) {
let current = env!("CARGO_PKG_VERSION");
let mut state = envelope::read_state(root);
if state.last_version.as_deref() == Some(current) {
return;
}
let kind = match state.last_version.as_deref() {
None => InstallKind::Install,
Some(previous) if version_is_older(previous, current) => InstallKind::Upgrade,
Some(_) => InstallKind::Downgrade,
};
let previous_version = state.last_version.clone();
state.schema_version = SCHEMA_VERSION;
state.last_version = Some(current.to_string());
if envelope::write_state(root, &state).is_err() {
return;
}
record(Event::InstallOrUpgrade {
kind,
previous_version,
});
}
fn version_is_older(previous: &str, current: &str) -> bool {
fn parts(value: &str) -> Vec<u64> {
value
.split(['-', '+'])
.next()
.unwrap_or_default()
.split('.')
.map(|part| part.parse::<u64>().unwrap_or_default())
.collect()
}
let (previous, current) = (parts(previous), parts(current));
let width = previous.len().max(current.len());
for index in 0..width {
let left = previous.get(index).copied().unwrap_or_default();
let right = current.get(index).copied().unwrap_or_default();
if left != right {
return left < right;
}
}
false
}
#[must_use]
pub fn is_armed() -> bool {
ARMED.get().is_some()
}
pub fn session_counters() -> &'static SessionCounters {
static COUNTERS: OnceLock<SessionCounters> = OnceLock::new();
COUNTERS.get_or_init(SessionCounters::default)
}
pub fn record(event: Event) {
let Some(armed) = ARMED.get() else {
return;
};
armed.handle.record(event);
}
pub fn record_blocking(event: Event) {
let Some(armed) = ARMED.get() else {
return;
};
let Ok(line) = serde_json::to_string(&event) else {
return;
};
let path = buffer::buffer_path(&armed.root);
let _ = buffer::append(&armed.root, &path, &line);
}
pub fn set_exit_class(class: ExitClass) {
let Some(armed) = ARMED.get() else {
return;
};
armed.exit_class.store(class.as_u8(), Ordering::Relaxed);
}
#[must_use]
pub fn exit_class() -> ExitClass {
ARMED.get().map_or(ExitClass::Clean, |armed| {
ExitClass::from_u8(armed.exit_class.load(Ordering::Relaxed))
})
}
pub fn shutdown_blocking(deadline: Duration) -> FlushOutcome {
ARMED
.get()
.map_or(FlushOutcome::Empty, |armed| armed.handle.shutdown(deadline))
}