mod policy;
mod upload;
use crate::wire::HookKind;
use crate::{Error, Result, host::GatewayHost};
use mobius::backend::model::provider::{HttpClient, HttpRedirectPolicy};
pub use policy::TelemetryPolicy;
use serde_json::{Value, json};
use std::sync::{
Arc, Mutex, RwLock,
atomic::{AtomicU64, Ordering},
};
use std::time::{Duration, Instant};
use serde::{Deserialize, Serialize};
use std::collections::BTreeMap;
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
#[serde(default, deny_unknown_fields)]
pub struct TelemetryConfig {
pub policy: TelemetryPolicy,
pub revision: u64,
pub sinks: Vec<TelemetrySink>,
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum SinkMethod {
#[default]
Post,
Get,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum TelemetrySection {
Activity,
Usage,
Runs,
Storage,
}
impl TelemetrySection {
const fn key(self) -> &'static str {
match self {
Self::Activity => "activity",
Self::Usage => "usage",
Self::Runs => "runs",
Self::Storage => "storage",
}
}
}
struct Snapshot {
clients: usize,
sections: BTreeMap<TelemetrySection, Value>,
}
impl Snapshot {
async fn read(&mut self, host: &GatewayHost, sections: &[TelemetrySection]) -> Result<Value> {
let mut result = host.telemetry_snapshot(&[], self.clients).await?;
for section in sections {
if !self.sections.contains_key(section) {
let mut snapshot = host
.telemetry_snapshot(std::slice::from_ref(section), self.clients)
.await?;
let value = snapshot
.as_object_mut()
.and_then(|object| object.remove(section.key()))
.ok_or_else(|| {
Error::Config(format!("telemetry {} section is missing", section.key()))
})?;
self.sections.insert(*section, value);
}
if let Some(value) = self.sections.get(section) {
result[section.key()] = value.clone();
}
}
Ok(result)
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub struct TelemetrySink {
pub id: String,
pub url: String,
#[serde(default)]
pub method: SinkMethod,
pub every_seconds: u32,
#[serde(default)]
pub sections: Vec<TelemetrySection>,
#[serde(default)]
pub events: Vec<HookKind>,
#[serde(default)]
pub headers: BTreeMap<String, String>,
#[serde(default)]
pub bearer_env: Option<String>,
#[serde(default)]
pub bearer_file: Option<String>,
#[serde(default)]
pub fields: BTreeMap<String, String>,
#[serde(default = "enabled")]
pub enabled: bool,
#[serde(default)]
pub upload_admission: bool,
}
impl TelemetrySink {
pub(crate) fn redact_report(&mut self) -> Result<()> {
let url = url::Url::parse(&self.url)
.map_err(|_| Error::Config("telemetry endpoint is invalid".into()))?;
self.url = url.origin().ascii_serialization();
self.headers.clear();
self.bearer_env = None;
self.bearer_file = None;
Ok(())
}
}
const fn enabled() -> bool {
true
}
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
pub struct TelemetrySinkStatus {
pub last_attempt_at: Option<i64>,
pub last_success_at: Option<i64>,
pub last_status: Option<u16>,
pub last_error: Option<String>,
pub consecutive_failures: u32,
pub next_at: Option<i64>,
pub events_pending: u64,
pub in_flight: bool,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct TelemetrySinkReport {
pub sink: TelemetrySink,
pub auth: SinkAuth,
pub status: TelemetrySinkStatus,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum SinkAuth {
None,
BearerEnv,
BearerFile,
}
#[derive(Debug, Clone, Copy, Serialize)]
#[serde(rename_all = "snake_case")]
pub(crate) enum StopCause {
Idle,
Signal,
LeaseExpired,
Error,
}
#[derive(Debug, Clone, Copy, Serialize)]
#[serde(tag = "reason", content = "cause", rename_all = "snake_case")]
pub(crate) enum Trigger {
Start,
Interval,
Events,
Manual,
Stop(StopCause),
}
pub(crate) struct Telemetry {
pub(crate) notify: Arc<tokio::sync::Notify>,
config: RwLock<Arc<TelemetryConfig>>,
state_dir: std::path::PathBuf,
client: Result<HttpClient>,
statuses: Mutex<BTreeMap<String, TelemetrySinkStatus>>,
manual: Mutex<std::collections::BTreeSet<String>>,
sequence: AtomicU64,
instance: String,
started_at_ms: i64,
started: Instant,
}
impl Telemetry {
pub(crate) fn new(config: &TelemetryConfig, state_dir: &std::path::Path) -> Self {
Self {
notify: Arc::new(tokio::sync::Notify::new()),
config: RwLock::new(Arc::new(config.clone())),
state_dir: state_dir.to_path_buf(),
client: HttpClient::builder()
.redirect(HttpRedirectPolicy::none())
.timeout(Duration::from_secs(config.policy.request_timeout_seconds))
.build()
.map_err(|error| {
Error::Config(format!(
"telemetry HTTP client initialization failed: {error}"
))
}),
statuses: Mutex::new(
config
.sinks
.iter()
.map(|sink| (sink.id.clone(), TelemetrySinkStatus::default()))
.collect(),
),
manual: Mutex::default(),
sequence: AtomicU64::new(0),
instance: uuid::Uuid::new_v4().to_string(),
started_at_ms: chrono::Utc::now().timestamp_millis(),
started: Instant::now(),
}
}
pub(crate) fn config(&self) -> Result<Arc<TelemetryConfig>> {
self.config
.read()
.map(|config| Arc::clone(&config))
.map_err(|_| Error::Config("telemetry configuration lock poisoned".into()))
}
fn header(&self, sink: &TelemetrySink, now: i64) -> Value {
json!({"version": 1, "sent_at": now, "sequence": self.sequence.fetch_add(1, Ordering::Relaxed),
"instance": self.instance, "gateway_version": env!("CARGO_PKG_VERSION"),
"protocol_version": crate::wire::PROTOCOL_VERSION, "started_at_ms": self.started_at_ms,
"uptime_seconds": self.started.elapsed().as_secs(), "fields": sink.fields})
}
pub(crate) fn configure(&self, config: TelemetryConfig) -> Result<()> {
let mut live = self
.config
.write()
.map_err(|_| Error::Config("telemetry configuration lock poisoned".into()))?;
let mut statuses = self
.statuses
.lock()
.map_err(|_| Error::Config("telemetry status lock poisoned".into()))?;
statuses.retain(|id, _| config.sinks.iter().any(|sink| &sink.id == id));
for sink in &config.sinks {
statuses.entry(sink.id.clone()).or_default();
}
for status in statuses.values_mut() {
status.next_at = None;
}
*live = Arc::new(config);
Ok(())
}
pub(crate) fn request_manual(&self, id: String) -> Result<()> {
if !self
.config()?
.sinks
.iter()
.any(|sink| sink.id == id && sink.enabled)
{
return Err(Error::Config("unknown or disabled telemetry sink".into()));
}
self.manual
.lock()
.map_err(|_| Error::Config("telemetry manual request lock poisoned".into()))?
.insert(id);
Ok(())
}
pub(crate) fn status(&self, id: &str) -> Result<TelemetrySinkStatus> {
Ok(self
.statuses
.lock()
.map_err(|_| Error::Config("telemetry status lock poisoned".into()))?
.get(id)
.cloned()
.unwrap_or_default())
}
pub(crate) fn can_drain(&self, id: &str) -> Result<bool> {
Ok(self
.statuses
.lock()
.map_err(|_| Error::Config("telemetry status lock poisoned".into()))?
.get(id)
.is_none_or(|status| status.consecutive_failures < 3))
}
pub(crate) async fn tick(
host: &GatewayHost,
clients: usize,
trigger: Trigger,
tasks: &mut tokio::task::JoinSet<()>,
) {
if !tasks.is_empty() {
return;
}
let host = host.clone();
tasks.spawn(async move {
let mut deliveries = tokio::task::JoinSet::new();
if let Err(error) = Self::tick_at(
&host,
clients,
trigger,
&mut deliveries,
chrono::Utc::now().timestamp(),
)
.await
{
eprintln!("telemetry scheduling failed: {error}");
}
while deliveries.join_next().await.is_some() {}
});
}
pub(crate) async fn stop(
host: &GatewayHost,
cause: StopCause,
tasks: &mut tokio::task::JoinSet<()>,
) {
tasks.shutdown().await;
match host.telemetry.statuses.lock() {
Ok(mut statuses) => {
for status in statuses.values_mut() {
status.in_flight = false;
}
}
Err(_) => eprintln!("telemetry stop status lock poisoned"),
}
Self::tick(host, 0, Trigger::Stop(cause), tasks).await;
if tokio::time::timeout(Duration::from_secs(5), async {
while tasks.join_next().await.is_some() {}
})
.await
.is_err()
{
eprintln!("telemetry stop delivery timed out");
}
tasks.shutdown().await;
}
pub(crate) async fn tick_at(
host: &GatewayHost,
clients: usize,
trigger: Trigger,
tasks: &mut tokio::task::JoinSet<()>,
now: i64,
) -> Result<()> {
let config = host.telemetry.config()?;
let mut snapshot = Snapshot {
clients,
sections: BTreeMap::new(),
};
for (index, sink) in config
.sinks
.iter()
.enumerate()
.filter(|(_, sink)| sink.enabled)
{
if let Err(error) = Self::schedule(
host,
trigger,
Arc::clone(&config),
index,
tasks,
now,
&mut snapshot,
)
.await
{
eprintln!("telemetry scheduling failed for {}: {error}", sink.id);
match host.telemetry.statuses.lock() {
Ok(mut statuses) => {
let Some(status) = statuses.get_mut(&sink.id) else {
continue;
};
status.in_flight = false;
status.last_attempt_at = Some(now);
status.last_error = Some(error.to_string());
status.consecutive_failures = status.consecutive_failures.saturating_add(1);
status.next_at = Some(now.saturating_add(i64::from(sink.every_seconds)));
}
Err(_) => eprintln!("telemetry status lock poisoned"),
}
}
}
Ok(())
}
async fn schedule(
host: &GatewayHost,
trigger: Trigger,
config: Arc<TelemetryConfig>,
index: usize,
tasks: &mut tokio::task::JoinSet<()>,
now: i64,
snapshot: &mut Snapshot,
) -> Result<()> {
let sink = &config.sinks[index];
let manual = host
.telemetry
.manual
.lock()
.map_err(|_| Error::Config("telemetry manual request lock poisoned".into()))?
.contains(&sink.id);
let (snapshot_due, retry_due) = {
let statuses = host
.telemetry
.statuses
.lock()
.map_err(|_| Error::Config("telemetry status lock poisoned".into()))?;
let status = statuses.get(&sink.id);
if status.is_some_and(|status| status.in_flight) {
return Ok(());
}
let snapshot_due = match trigger {
Trigger::Events => false,
Trigger::Interval => {
manual
|| status
.and_then(|status| status.next_at)
.is_none_or(|at| at <= now)
}
Trigger::Start | Trigger::Manual | Trigger::Stop(_) => true,
};
let retry_due = status.is_none_or(|status| {
status.consecutive_failures == 0
|| status
.last_attempt_at
.is_none_or(|at| now.saturating_sub(at) >= 15)
});
(snapshot_due, retry_due)
};
let mut cursor = None;
let mut pending = 0;
let mut envelope = if snapshot_due {
let sections = if matches!(trigger, Trigger::Stop(_)) {
&[][..]
} else {
&sink.sections
};
snapshot.read(host, sections).await?
} else {
if sink.events.is_empty() || !retry_due {
return Ok(());
}
let (events, after, count) = host.telemetry_events(sink).await?;
if events.is_empty() {
return Ok(());
}
cursor = after;
pending = count;
let mut header = snapshot.read(host, &[]).await?;
header["events"] = json!(events);
header
};
let reason = if !snapshot_due {
Trigger::Events
} else if manual && matches!(trigger, Trigger::Interval) {
Trigger::Manual
} else {
trigger
};
let header = host.telemetry.header(sink, now);
let object = envelope
.as_object_mut()
.ok_or_else(|| Error::Config("telemetry snapshot is not an object".into()))?;
if let Value::Object(header) = header {
object.extend(header);
}
if let Value::Object(reason) = serde_json::to_value(reason)? {
object.extend(reason);
}
{
let mut statuses = host
.telemetry
.statuses
.lock()
.map_err(|_| Error::Config("telemetry status lock poisoned".into()))?;
let Some(status) = statuses.get_mut(&sink.id) else {
return Ok(());
};
status.in_flight = true;
status.last_attempt_at = Some(now);
status.events_pending = pending;
if snapshot_due {
status.next_at = Some(now.saturating_add(i64::from(sink.every_seconds)));
}
}
if manual {
host.telemetry
.manual
.lock()
.map_err(|_| Error::Config("telemetry manual request lock poisoned".into()))?
.remove(&sink.id);
}
let host = host.clone();
tasks.spawn(async move {
let sink = &config.sinks[index];
let result = deliver(&host.telemetry, sink, &envelope, 0).await;
let http = match &result {
Ok((status, _)) => Some(*status),
Err(error) => error.status,
};
let acknowledge = match &result {
Ok(_) => true,
Err(error) => error.permanent,
};
let mut error = result.err().map(|error| error.message);
if acknowledge
&& let Some(cursor) = cursor
&& let Err(failure) = host
.advance_telemetry(&sink.id, cursor, config.revision)
.await
{
error = Some(failure.to_string());
}
match host.telemetry.statuses.lock() {
Ok(mut statuses) => {
if let Some(status) = statuses.get_mut(&sink.id) {
status.in_flight = false;
status.last_status = http;
if error.is_none() {
status.last_success_at = Some(chrono::Utc::now().timestamp());
status.consecutive_failures = 0;
} else {
status.consecutive_failures =
status.consecutive_failures.saturating_add(1);
}
status.last_error = error;
}
}
Err(_) => eprintln!("telemetry delivery status lock poisoned"),
}
host.telemetry.notify.notify_one();
});
Ok(())
}
pub(crate) async fn pending(host: &GatewayHost) -> bool {
match host.telemetry_pending().await {
Ok(pending) => pending,
Err(error) => {
eprintln!("telemetry pending check failed: {error}");
false
}
}
}
}
#[derive(Debug, thiserror::Error)]
#[error("{message}")]
struct DeliveryError {
status: Option<u16>,
message: String,
permanent: bool,
}
impl DeliveryError {
fn caused(context: &str, cause: &dyn std::error::Error) -> Self {
use std::fmt::Write as _;
let mut message = format!("{context}: {cause}");
let mut source = cause.source();
while let Some(cause) = source {
let _ = write!(message, ": {cause}");
source = cause.source();
}
Self {
status: None,
message,
permanent: false,
}
}
fn transient(message: &str) -> Self {
Self {
status: None,
message: message.into(),
permanent: false,
}
}
}
async fn deliver(
telemetry: &Telemetry,
sink: &TelemetrySink,
envelope: &Value,
response_limit: usize,
) -> std::result::Result<(u16, Vec<u8>), DeliveryError> {
let error = DeliveryError::transient;
let client = telemetry.client.as_ref().map_err(|cause| {
error(&format!(
"telemetry HTTP client initialization failed: {cause}"
))
})?;
let mut request = match sink.method {
SinkMethod::Post => {
let bytes = serde_json::to_vec(envelope)
.map_err(|cause| error(&format!("telemetry encoding failed: {cause}")))?;
if bytes.len() > 64 * 1024 {
return Err(DeliveryError {
status: None,
message: "telemetry envelope exceeds 64 KiB".into(),
permanent: true,
});
}
client
.post(&sink.url)
.header("content-type", "application/json")
.body(bytes)
}
SinkMethod::Get => client.get(&sink.url),
};
for (name, value) in &sink.headers {
request = request.header(name, value);
}
let token = if let Some(name) = &sink.bearer_env {
Some(std::env::var(name).map_err(|_| error("bearer environment variable unavailable"))?)
} else if let Some(path) = &sink.bearer_file {
let path = telemetry.state_dir.join(path);
let state_dir = &telemetry.state_dir;
let canonical = tokio::fs::canonicalize(&path)
.await
.map_err(|cause| error(&format!("bearer file unavailable: {cause}")))?;
if !canonical.starts_with(state_dir) {
return Err(error("bearer file escapes state directory"));
}
Some(
tokio::task::spawn_blocking(move || crate::config::load_secret_file(&path))
.await
.map_err(|cause| error(&format!("bearer file task failed: {cause}")))?
.map_err(|cause| error(&format!("invalid bearer file: {cause}")))?,
)
} else {
None
};
if let Some(token) = token {
if token.is_empty()
|| token.len() > 16 * 1024
|| !token.bytes().all(|b| (33..=126).contains(&b))
{
return Err(error("invalid bearer token"));
}
request = request.bearer_auth(token);
}
let mut response = request
.send()
.await
.map_err(|cause| DeliveryError::caused("telemetry request failed", &cause.without_url()))?;
let status = response.status();
if status.is_success() {
let mut body = Vec::new();
if response_limit > 0 {
while let Some(chunk) = response.chunk().await.map_err(|cause| {
DeliveryError::caused("telemetry response failed", &cause.without_url())
})? {
if chunk.len() > response_limit.saturating_sub(body.len()) {
return Err(error("telemetry response exceeds its size limit"));
}
body.extend_from_slice(&chunk);
}
}
return Ok((status.as_u16(), body));
}
Err(DeliveryError {
status: Some(status.as_u16()),
message: format!("telemetry collector returned HTTP {}", status.as_u16()),
permanent: status.is_client_error() && status.as_u16() != 408 && status.as_u16() != 429,
})
}