1mod activity;
3mod policy;
4mod resources;
5mod upload;
6use crate::wire::HookKind;
7use crate::{Error, Result, host::GatewayHost};
8pub(crate) use activity::ActivityHook;
9pub use activity::ActivityHookConfig;
10use mobius::backend::model::provider::{HttpClient, HttpRedirectPolicy};
11pub use policy::TelemetryPolicy;
12use serde_json::{Value, json};
13use std::sync::{
14 Arc, Mutex, RwLock,
15 atomic::{AtomicU64, Ordering},
16};
17use std::time::{Duration, Instant};
18
19use serde::{Deserialize, Serialize};
20use std::collections::BTreeMap;
21
22#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
24#[serde(default, deny_unknown_fields)]
25pub struct TelemetryConfig {
26 pub policy: TelemetryPolicy,
28 #[serde(skip_serializing_if = "Option::is_none")]
30 pub activity_hook: Option<ActivityHookConfig>,
31 pub revision: u64,
33 pub sinks: Vec<TelemetrySink>,
35}
36#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
38#[serde(rename_all = "snake_case")]
39pub enum SinkMethod {
40 #[default]
42 Post,
43 Get,
45}
46#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize)]
48#[serde(rename_all = "snake_case")]
49pub enum TelemetrySection {
50 Activity,
52 Usage,
54 Runs,
56 Storage,
58 Resources,
60}
61impl TelemetrySection {
62 const fn key(self) -> &'static str {
63 match self {
64 Self::Activity => "activity",
65 Self::Usage => "usage",
66 Self::Runs => "runs",
67 Self::Storage => "storage",
68 Self::Resources => "resources",
69 }
70 }
71}
72
73struct Snapshot {
75 clients: usize,
76 sections: BTreeMap<TelemetrySection, Value>,
77}
78impl Snapshot {
79 async fn read(&mut self, host: &GatewayHost, sections: &[TelemetrySection]) -> Result<Value> {
80 let mut result = host.telemetry_snapshot(&[], self.clients).await?;
81 for section in sections {
82 if !self.sections.contains_key(section) {
83 let mut snapshot = host
84 .telemetry_snapshot(std::slice::from_ref(section), self.clients)
85 .await?;
86 let value = snapshot
87 .as_object_mut()
88 .and_then(|object| object.remove(section.key()))
89 .ok_or_else(|| {
90 Error::Config(format!("telemetry {} section is missing", section.key()))
91 })?;
92 self.sections.insert(*section, value);
93 }
94 if let Some(value) = self.sections.get(section) {
95 result[section.key()] = value.clone();
97 }
98 }
99 Ok(result)
100 }
101}
102
103#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
105#[serde(deny_unknown_fields)]
106pub struct TelemetrySink {
107 pub id: String,
109 pub url: String,
111 #[serde(default)]
113 pub method: SinkMethod,
114 pub every_seconds: u32,
116 #[serde(default)]
118 pub sections: Vec<TelemetrySection>,
119 #[serde(default)]
121 pub events: Vec<HookKind>,
122 #[serde(default)]
124 pub headers: BTreeMap<String, String>,
125 #[serde(default)]
127 pub bearer_env: Option<String>,
128 #[serde(default)]
130 pub bearer_file: Option<String>,
131 #[serde(default)]
133 pub fields: BTreeMap<String, String>,
134 #[serde(default = "enabled")]
136 pub enabled: bool,
137 #[serde(default)]
139 pub upload_admission: bool,
140}
141impl TelemetrySink {
142 pub(crate) fn redact_report(&mut self) -> Result<()> {
143 let url = url::Url::parse(&self.url)
144 .map_err(|_| Error::Config("telemetry endpoint is invalid".into()))?;
145 self.url = url.origin().ascii_serialization();
147 self.headers.clear();
148 self.bearer_env = None;
149 self.bearer_file = None;
150 Ok(())
151 }
152}
153
154const fn enabled() -> bool {
155 true
156}
157#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
159pub struct TelemetrySinkStatus {
160 pub last_attempt_at: Option<i64>,
162 pub last_success_at: Option<i64>,
164 pub last_status: Option<u16>,
166 pub last_error: Option<String>,
168 pub consecutive_failures: u32,
170 pub next_at: Option<i64>,
172 pub events_pending: u64,
174 pub in_flight: bool,
176}
177#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
179pub struct TelemetrySinkReport {
180 pub sink: TelemetrySink,
182 pub auth: SinkAuth,
184 pub status: TelemetrySinkStatus,
186}
187
188#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
190#[serde(rename_all = "snake_case")]
191pub enum SinkAuth {
192 None,
194 BearerEnv,
196 BearerFile,
198}
199
200#[derive(Debug, Clone, Copy, Serialize)]
201#[serde(rename_all = "snake_case")]
202pub(crate) enum StopCause {
203 Idle,
204 Signal,
205 LeaseExpired,
206 Error,
207}
208
209#[derive(Debug, Clone, Copy, Serialize)]
210#[serde(tag = "reason", content = "cause", rename_all = "snake_case")]
211pub(crate) enum Trigger {
212 Start,
213 Interval,
214 Events,
215 Manual,
216 Stop(StopCause),
217}
218
219pub(crate) struct Telemetry {
220 pub(crate) notify: Arc<tokio::sync::Notify>,
221 config: RwLock<Arc<TelemetryConfig>>,
224 state_dir: std::path::PathBuf,
225 client: Result<HttpClient>,
226 statuses: Mutex<BTreeMap<String, TelemetrySinkStatus>>,
227 manual: Mutex<std::collections::BTreeSet<String>>,
228 sequence: AtomicU64,
229 instance: String,
230 started_at_ms: i64,
231 started: Instant,
232 resources: resources::Resources,
233}
234impl Telemetry {
235 pub(crate) fn new(config: &TelemetryConfig, state_dir: &std::path::Path) -> Self {
236 Self {
237 notify: Arc::new(tokio::sync::Notify::new()),
238 config: RwLock::new(Arc::new(config.clone())),
239 state_dir: state_dir.to_path_buf(),
240 client: HttpClient::builder()
242 .redirect(HttpRedirectPolicy::none())
243 .timeout(Duration::from_secs(config.policy.request_timeout_seconds))
244 .build()
245 .map_err(|error| {
246 Error::Config(format!(
247 "telemetry HTTP client initialization failed: {error}"
248 ))
249 }),
250 statuses: Mutex::new(
251 config
252 .sinks
253 .iter()
254 .map(|sink| (sink.id.clone(), TelemetrySinkStatus::default()))
255 .collect(),
256 ),
257 manual: Mutex::default(),
258 sequence: AtomicU64::new(0),
259 instance: uuid::Uuid::new_v4().to_string(),
260 started_at_ms: chrono::Utc::now().timestamp_millis(),
261 started: Instant::now(),
262 resources: resources::Resources::default(),
263 }
264 }
265 pub(crate) async fn resources(&self) -> Value {
266 self.resources.sample().await
267 }
268
269 pub(crate) fn config(&self) -> Result<Arc<TelemetryConfig>> {
270 self.config
271 .read()
272 .map(|config| Arc::clone(&config))
273 .map_err(|_| Error::Config("telemetry configuration lock poisoned".into()))
274 }
275
276 fn header(&self, sink: &TelemetrySink, now: i64) -> Value {
277 json!({"version": 1, "sent_at": now, "sequence": self.sequence.fetch_add(1, Ordering::Relaxed),
278 "instance": self.instance, "gateway_version": env!("CARGO_PKG_VERSION"),
279 "protocol_version": crate::wire::PROTOCOL_VERSION, "started_at_ms": self.started_at_ms,
280 "uptime_seconds": self.started.elapsed().as_secs(), "fields": sink.fields})
281 }
282 pub(crate) fn configure_after(
284 &self,
285 operation: impl FnOnce() -> Result<(TelemetryConfig, crate::publication::Outcome)>,
286 ) -> Result<()> {
287 let mut live = self
288 .config
289 .write()
290 .map_err(|_| Error::Config("telemetry configuration lock poisoned".into()))?;
291 let mut statuses = self
292 .statuses
293 .lock()
294 .map_err(|_| Error::Config("telemetry status lock poisoned".into()))?;
295 let (config, publication) = operation()?;
296 statuses.retain(|id, _| config.sinks.iter().any(|sink| &sink.id == id));
297 for sink in &config.sinks {
298 match statuses.get_mut(&sink.id) {
299 Some(status) => status.next_at = None,
300 None => {
301 statuses.insert(sink.id.clone(), TelemetrySinkStatus::default());
302 }
303 }
304 }
305 *live = Arc::new(config);
306 publication.confirm()
307 }
308 pub(crate) fn request_manual(&self, id: String) -> Result<()> {
309 if !self
310 .config()?
311 .sinks
312 .iter()
313 .any(|sink| sink.id == id && sink.enabled)
314 {
315 return Err(Error::Config("unknown or disabled telemetry sink".into()));
316 }
317 self.manual
318 .lock()
319 .map_err(|_| Error::Config("telemetry manual request lock poisoned".into()))?
320 .insert(id);
321 Ok(())
322 }
323 pub(crate) fn status(&self, id: &str) -> Result<TelemetrySinkStatus> {
324 Ok(self
326 .statuses
327 .lock()
328 .map_err(|_| Error::Config("telemetry status lock poisoned".into()))?
329 .get(id)
330 .cloned()
331 .unwrap_or_default())
332 }
333 pub(crate) fn can_drain(&self, id: &str) -> Result<bool> {
334 Ok(self
335 .statuses
336 .lock()
337 .map_err(|_| Error::Config("telemetry status lock poisoned".into()))?
338 .get(id)
339 .is_none_or(|status| status.consecutive_failures < 3))
340 }
341 pub(crate) async fn tick(
343 host: &GatewayHost,
344 clients: usize,
345 trigger: Trigger,
346 tasks: &mut tokio::task::JoinSet<()>,
347 ) {
348 if !tasks.is_empty() {
349 return;
350 }
351 let config = match host.telemetry.config() {
352 Ok(config) => config,
353 Err(error) => {
354 tracing::warn!(%error, "telemetry scheduling failed");
355 return;
356 }
357 };
358 if !config.sinks.iter().any(|sink| sink.enabled) {
359 return;
360 }
361 let host = host.clone();
362 tasks.spawn(async move {
363 let mut deliveries = tokio::task::JoinSet::new();
364 if let Err(error) = Self::tick_at(
365 &host,
366 clients,
367 trigger,
368 &mut deliveries,
369 chrono::Utc::now().timestamp(),
370 )
371 .await
372 {
373 tracing::warn!(%error, "telemetry scheduling failed");
374 }
375 while deliveries.join_next().await.is_some() {}
376 });
377 }
378 pub(crate) async fn stop(
379 host: &GatewayHost,
380 cause: StopCause,
381 tasks: &mut tokio::task::JoinSet<()>,
382 ) {
383 tasks.shutdown().await;
386 match host.telemetry.statuses.lock() {
387 Ok(mut statuses) => {
388 for status in statuses.values_mut() {
389 status.in_flight = false;
390 }
391 }
392 Err(_) => tracing::warn!("telemetry stop status lock poisoned"),
393 }
394 Self::tick(host, 0, Trigger::Stop(cause), tasks).await;
395 if tokio::time::timeout(Duration::from_secs(5), async {
396 while tasks.join_next().await.is_some() {}
397 })
398 .await
399 .is_err()
400 {
401 tracing::warn!("telemetry stop delivery timed out");
402 }
403 tasks.shutdown().await;
404 }
405 pub(crate) async fn tick_at(
406 host: &GatewayHost,
407 clients: usize,
408 trigger: Trigger,
409 tasks: &mut tokio::task::JoinSet<()>,
410 now: i64,
411 ) -> Result<()> {
412 let config = host.telemetry.config()?;
413 let mut snapshot = Snapshot {
414 clients,
415 sections: BTreeMap::new(),
416 };
417 for (index, sink) in config
418 .sinks
419 .iter()
420 .enumerate()
421 .filter(|(_, sink)| sink.enabled)
422 {
423 if let Err(error) = Self::schedule(
425 host,
426 trigger,
427 Arc::clone(&config),
428 index,
429 tasks,
430 now,
431 &mut snapshot,
432 )
433 .await
434 {
435 tracing::warn!(sink_id = %sink.id, %error, "telemetry scheduling failed");
436 match host.telemetry.statuses.lock() {
437 Ok(mut statuses) => {
438 let Some(status) = statuses.get_mut(&sink.id) else {
439 continue;
440 };
441 status.in_flight = false;
442 status.last_attempt_at = Some(now);
443 status.last_error = Some(error.to_string());
444 status.consecutive_failures = status.consecutive_failures.saturating_add(1);
445 status.next_at = Some(now.saturating_add(i64::from(sink.every_seconds)));
446 }
447 Err(_) => tracing::warn!("telemetry status lock poisoned"),
448 }
449 }
450 }
451 Ok(())
452 }
453 async fn schedule(
454 host: &GatewayHost,
455 trigger: Trigger,
456 config: Arc<TelemetryConfig>,
457 index: usize,
458 tasks: &mut tokio::task::JoinSet<()>,
459 now: i64,
460 snapshot: &mut Snapshot,
461 ) -> Result<()> {
462 let sink = &config.sinks[index];
463 let manual = host
464 .telemetry
465 .manual
466 .lock()
467 .map_err(|_| Error::Config("telemetry manual request lock poisoned".into()))?
468 .contains(&sink.id);
469 let (snapshot_due, retry_due) = {
470 let statuses = host
471 .telemetry
472 .statuses
473 .lock()
474 .map_err(|_| Error::Config("telemetry status lock poisoned".into()))?;
475 let status = statuses.get(&sink.id);
476 if status.is_some_and(|status| status.in_flight) {
477 return Ok(());
478 }
479 let snapshot_due = match trigger {
480 Trigger::Events => false,
481 Trigger::Interval => {
482 manual
483 || status
484 .and_then(|status| status.next_at)
485 .is_none_or(|at| at <= now)
486 }
487 Trigger::Start | Trigger::Manual | Trigger::Stop(_) => true,
488 };
489 let retry_due = status.is_none_or(|status| {
490 status.consecutive_failures == 0
491 || status
492 .last_attempt_at
493 .is_none_or(|at| now.saturating_sub(at) >= 15)
494 });
495 (snapshot_due, retry_due)
496 };
497 let mut cursor = None;
498 let mut pending = 0;
499 let mut envelope = if snapshot_due {
500 let sections = if matches!(trigger, Trigger::Stop(_)) {
501 &[][..]
502 } else {
503 &sink.sections
504 };
505 snapshot.read(host, sections).await?
506 } else {
507 if sink.events.is_empty() || !retry_due {
508 return Ok(());
509 }
510 let (events, after, count) = host.telemetry_events(sink).await?;
511 if events.is_empty() {
512 return Ok(());
513 }
514 cursor = after;
515 pending = count;
516 let mut header = snapshot.read(host, &[]).await?;
517 header["events"] = json!(events);
518 header
519 };
520 let reason = if !snapshot_due {
521 Trigger::Events
522 } else if manual && matches!(trigger, Trigger::Interval) {
523 Trigger::Manual
524 } else {
525 trigger
526 };
527 let header = host.telemetry.header(sink, now);
528 let object = envelope
529 .as_object_mut()
530 .ok_or_else(|| Error::Config("telemetry snapshot is not an object".into()))?;
531 if let Value::Object(header) = header {
532 object.extend(header);
533 }
534 if let Value::Object(reason) = serde_json::to_value(reason)? {
535 object.extend(reason);
536 }
537 {
538 let mut statuses = host
539 .telemetry
540 .statuses
541 .lock()
542 .map_err(|_| Error::Config("telemetry status lock poisoned".into()))?;
543 let Some(status) = statuses.get_mut(&sink.id) else {
544 return Ok(());
545 };
546 status.in_flight = true;
547 status.last_attempt_at = Some(now);
548 status.events_pending = pending;
549 if snapshot_due {
550 status.next_at = Some(now.saturating_add(i64::from(sink.every_seconds)));
551 }
552 }
553 if manual {
554 host.telemetry
555 .manual
556 .lock()
557 .map_err(|_| Error::Config("telemetry manual request lock poisoned".into()))?
558 .remove(&sink.id);
559 }
560 let host = host.clone();
561 tasks.spawn(async move {
562 let sink = &config.sinks[index];
563 let result = deliver(&host.telemetry, sink, &envelope, 0).await;
564 let http = match &result {
565 Ok((status, _)) => Some(*status),
566 Err(error) => error.status,
567 };
568 let acknowledge = match &result {
569 Ok(_) => true,
570 Err(error) => error.permanent,
571 };
572 let mut error = result.err().map(|error| error.message);
573 if acknowledge
574 && let Some(cursor) = cursor
575 && let Err(failure) = host
576 .advance_telemetry(&sink.id, cursor, config.revision)
577 .await
578 {
579 error = Some(failure.to_string());
580 }
581 match host.telemetry.statuses.lock() {
582 Ok(mut statuses) => {
583 if let Some(status) = statuses.get_mut(&sink.id) {
584 status.in_flight = false;
585 status.last_status = http;
586 if error.is_none() {
587 status.last_success_at = Some(chrono::Utc::now().timestamp());
588 status.consecutive_failures = 0;
589 } else {
590 status.consecutive_failures =
591 status.consecutive_failures.saturating_add(1);
592 }
593 status.last_error = error;
594 }
595 }
596 Err(_) => tracing::warn!("telemetry delivery status lock poisoned"),
597 }
598 host.telemetry.notify.notify_one();
599 });
600 Ok(())
601 }
602 pub(crate) async fn pending(host: &GatewayHost) -> bool {
603 match host.telemetry_pending().await {
604 Ok(pending) => pending,
605 Err(error) => {
606 tracing::warn!(%error, "telemetry pending check failed");
607 false
608 }
609 }
610 }
611}
612
613#[derive(Debug, thiserror::Error)]
614#[error("{message}")]
615struct DeliveryError {
616 status: Option<u16>,
617 message: String,
618 permanent: bool,
619}
620impl DeliveryError {
621 fn caused(context: &str, cause: &dyn std::error::Error) -> Self {
622 use std::fmt::Write as _;
623 let mut message = format!("{context}: {cause}");
624 let mut source = cause.source();
625 while let Some(cause) = source {
626 let _ = write!(message, ": {cause}");
628 source = cause.source();
629 }
630 Self {
631 status: None,
632 message,
633 permanent: false,
634 }
635 }
636 fn transient(message: &str) -> Self {
637 Self {
638 status: None,
639 message: message.into(),
640 permanent: false,
641 }
642 }
643}
644async fn deliver(
645 telemetry: &Telemetry,
646 sink: &TelemetrySink,
647 envelope: &Value,
648 response_limit: usize,
649) -> std::result::Result<(u16, Vec<u8>), DeliveryError> {
650 let error = DeliveryError::transient;
651 let client = telemetry.client.as_ref().map_err(|cause| {
652 error(&format!(
653 "telemetry HTTP client initialization failed: {cause}"
654 ))
655 })?;
656 let mut request = match sink.method {
657 SinkMethod::Post => {
658 let bytes = serde_json::to_vec(envelope)
659 .map_err(|cause| error(&format!("telemetry encoding failed: {cause}")))?;
660 if bytes.len() > 64 * 1024 {
661 return Err(DeliveryError {
662 status: None,
663 message: "telemetry envelope exceeds 64 KiB".into(),
664 permanent: true,
665 });
666 }
667 client
668 .post(&sink.url)
669 .header("content-type", "application/json")
670 .body(bytes)
671 }
672 SinkMethod::Get => client.get(&sink.url),
673 };
674 for (name, value) in &sink.headers {
675 request = request.header(name, value);
676 }
677 let token = if let Some(name) = &sink.bearer_env {
678 Some(std::env::var(name).map_err(|_| error("bearer environment variable unavailable"))?)
679 } else if let Some(path) = &sink.bearer_file {
680 let path = telemetry.state_dir.join(path);
681 let state_dir = &telemetry.state_dir;
682 let canonical = tokio::fs::canonicalize(&path)
683 .await
684 .map_err(|cause| error(&format!("bearer file unavailable: {cause}")))?;
685 if !canonical.starts_with(state_dir) {
686 return Err(error("bearer file escapes state directory"));
687 }
688 Some(
689 tokio::task::spawn_blocking(move || crate::config::load_secret_file(&path))
690 .await
691 .map_err(|cause| error(&format!("bearer file task failed: {cause}")))?
692 .map_err(|cause| error(&format!("invalid bearer file: {cause}")))?,
693 )
694 } else {
695 None
696 };
697 if let Some(token) = token {
698 if token.is_empty()
699 || token.len() > 16 * 1024
700 || !token.bytes().all(|b| (33..=126).contains(&b))
701 {
702 return Err(error("invalid bearer token"));
703 }
704 request = request.bearer_auth(token);
705 }
706 let mut response = request
708 .send()
709 .await
710 .map_err(|cause| DeliveryError::caused("telemetry request failed", &cause.without_url()))?;
711 let status = response.status();
712 if status.is_success() {
713 let mut body = Vec::new();
714 if response_limit > 0 {
715 while let Some(chunk) = response.chunk().await.map_err(|cause| {
716 DeliveryError::caused("telemetry response failed", &cause.without_url())
717 })? {
718 if chunk.len() > response_limit.saturating_sub(body.len()) {
719 return Err(error("telemetry response exceeds its size limit"));
720 }
721 body.extend_from_slice(&chunk);
722 }
723 }
724 return Ok((status.as_u16(), body));
725 }
726 Err(DeliveryError {
727 status: Some(status.as_u16()),
728 message: format!("telemetry collector returned HTTP {}", status.as_u16()),
729 permanent: status.is_client_error() && status.as_u16() != 408 && status.as_u16() != 429,
730 })
731}