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(&self, config: TelemetryConfig) -> Result<()> {
283 let mut live = self
284 .config
285 .write()
286 .map_err(|_| Error::Config("telemetry configuration lock poisoned".into()))?;
287 let mut statuses = self
288 .statuses
289 .lock()
290 .map_err(|_| Error::Config("telemetry status lock poisoned".into()))?;
291 statuses.retain(|id, _| config.sinks.iter().any(|sink| &sink.id == id));
292 for sink in &config.sinks {
293 statuses.entry(sink.id.clone()).or_default();
294 }
295 for status in statuses.values_mut() {
296 status.next_at = None;
297 }
298 *live = Arc::new(config);
299 Ok(())
300 }
301 pub(crate) fn request_manual(&self, id: String) -> Result<()> {
302 if !self
303 .config()?
304 .sinks
305 .iter()
306 .any(|sink| sink.id == id && sink.enabled)
307 {
308 return Err(Error::Config("unknown or disabled telemetry sink".into()));
309 }
310 self.manual
311 .lock()
312 .map_err(|_| Error::Config("telemetry manual request lock poisoned".into()))?
313 .insert(id);
314 Ok(())
315 }
316 pub(crate) fn status(&self, id: &str) -> Result<TelemetrySinkStatus> {
317 Ok(self
319 .statuses
320 .lock()
321 .map_err(|_| Error::Config("telemetry status lock poisoned".into()))?
322 .get(id)
323 .cloned()
324 .unwrap_or_default())
325 }
326 pub(crate) fn can_drain(&self, id: &str) -> Result<bool> {
327 Ok(self
328 .statuses
329 .lock()
330 .map_err(|_| Error::Config("telemetry status lock poisoned".into()))?
331 .get(id)
332 .is_none_or(|status| status.consecutive_failures < 3))
333 }
334 pub(crate) async fn tick(
336 host: &GatewayHost,
337 clients: usize,
338 trigger: Trigger,
339 tasks: &mut tokio::task::JoinSet<()>,
340 ) {
341 if !tasks.is_empty() {
342 return;
343 }
344 let config = match host.telemetry.config() {
345 Ok(config) => config,
346 Err(error) => {
347 eprintln!("telemetry scheduling failed: {error}");
348 return;
349 }
350 };
351 if !config.sinks.iter().any(|sink| sink.enabled) {
352 return;
353 }
354 let host = host.clone();
355 tasks.spawn(async move {
356 let mut deliveries = tokio::task::JoinSet::new();
357 if let Err(error) = Self::tick_at(
358 &host,
359 clients,
360 trigger,
361 &mut deliveries,
362 chrono::Utc::now().timestamp(),
363 )
364 .await
365 {
366 eprintln!("telemetry scheduling failed: {error}");
367 }
368 while deliveries.join_next().await.is_some() {}
369 });
370 }
371 pub(crate) async fn stop(
372 host: &GatewayHost,
373 cause: StopCause,
374 tasks: &mut tokio::task::JoinSet<()>,
375 ) {
376 tasks.shutdown().await;
379 match host.telemetry.statuses.lock() {
380 Ok(mut statuses) => {
381 for status in statuses.values_mut() {
382 status.in_flight = false;
383 }
384 }
385 Err(_) => eprintln!("telemetry stop status lock poisoned"),
386 }
387 Self::tick(host, 0, Trigger::Stop(cause), tasks).await;
388 if tokio::time::timeout(Duration::from_secs(5), async {
389 while tasks.join_next().await.is_some() {}
390 })
391 .await
392 .is_err()
393 {
394 eprintln!("telemetry stop delivery timed out");
395 }
396 tasks.shutdown().await;
397 }
398 pub(crate) async fn tick_at(
399 host: &GatewayHost,
400 clients: usize,
401 trigger: Trigger,
402 tasks: &mut tokio::task::JoinSet<()>,
403 now: i64,
404 ) -> Result<()> {
405 let config = host.telemetry.config()?;
406 let mut snapshot = Snapshot {
407 clients,
408 sections: BTreeMap::new(),
409 };
410 for (index, sink) in config
411 .sinks
412 .iter()
413 .enumerate()
414 .filter(|(_, sink)| sink.enabled)
415 {
416 if let Err(error) = Self::schedule(
418 host,
419 trigger,
420 Arc::clone(&config),
421 index,
422 tasks,
423 now,
424 &mut snapshot,
425 )
426 .await
427 {
428 eprintln!("telemetry scheduling failed for {}: {error}", sink.id);
429 match host.telemetry.statuses.lock() {
430 Ok(mut statuses) => {
431 let Some(status) = statuses.get_mut(&sink.id) else {
432 continue;
433 };
434 status.in_flight = false;
435 status.last_attempt_at = Some(now);
436 status.last_error = Some(error.to_string());
437 status.consecutive_failures = status.consecutive_failures.saturating_add(1);
438 status.next_at = Some(now.saturating_add(i64::from(sink.every_seconds)));
439 }
440 Err(_) => eprintln!("telemetry status lock poisoned"),
441 }
442 }
443 }
444 Ok(())
445 }
446 async fn schedule(
447 host: &GatewayHost,
448 trigger: Trigger,
449 config: Arc<TelemetryConfig>,
450 index: usize,
451 tasks: &mut tokio::task::JoinSet<()>,
452 now: i64,
453 snapshot: &mut Snapshot,
454 ) -> Result<()> {
455 let sink = &config.sinks[index];
456 let manual = host
457 .telemetry
458 .manual
459 .lock()
460 .map_err(|_| Error::Config("telemetry manual request lock poisoned".into()))?
461 .contains(&sink.id);
462 let (snapshot_due, retry_due) = {
463 let statuses = host
464 .telemetry
465 .statuses
466 .lock()
467 .map_err(|_| Error::Config("telemetry status lock poisoned".into()))?;
468 let status = statuses.get(&sink.id);
469 if status.is_some_and(|status| status.in_flight) {
470 return Ok(());
471 }
472 let snapshot_due = match trigger {
473 Trigger::Events => false,
474 Trigger::Interval => {
475 manual
476 || status
477 .and_then(|status| status.next_at)
478 .is_none_or(|at| at <= now)
479 }
480 Trigger::Start | Trigger::Manual | Trigger::Stop(_) => true,
481 };
482 let retry_due = status.is_none_or(|status| {
483 status.consecutive_failures == 0
484 || status
485 .last_attempt_at
486 .is_none_or(|at| now.saturating_sub(at) >= 15)
487 });
488 (snapshot_due, retry_due)
489 };
490 let mut cursor = None;
491 let mut pending = 0;
492 let mut envelope = if snapshot_due {
493 let sections = if matches!(trigger, Trigger::Stop(_)) {
494 &[][..]
495 } else {
496 &sink.sections
497 };
498 snapshot.read(host, sections).await?
499 } else {
500 if sink.events.is_empty() || !retry_due {
501 return Ok(());
502 }
503 let (events, after, count) = host.telemetry_events(sink).await?;
504 if events.is_empty() {
505 return Ok(());
506 }
507 cursor = after;
508 pending = count;
509 let mut header = snapshot.read(host, &[]).await?;
510 header["events"] = json!(events);
511 header
512 };
513 let reason = if !snapshot_due {
514 Trigger::Events
515 } else if manual && matches!(trigger, Trigger::Interval) {
516 Trigger::Manual
517 } else {
518 trigger
519 };
520 let header = host.telemetry.header(sink, now);
521 let object = envelope
522 .as_object_mut()
523 .ok_or_else(|| Error::Config("telemetry snapshot is not an object".into()))?;
524 if let Value::Object(header) = header {
525 object.extend(header);
526 }
527 if let Value::Object(reason) = serde_json::to_value(reason)? {
528 object.extend(reason);
529 }
530 {
531 let mut statuses = host
532 .telemetry
533 .statuses
534 .lock()
535 .map_err(|_| Error::Config("telemetry status lock poisoned".into()))?;
536 let Some(status) = statuses.get_mut(&sink.id) else {
537 return Ok(());
538 };
539 status.in_flight = true;
540 status.last_attempt_at = Some(now);
541 status.events_pending = pending;
542 if snapshot_due {
543 status.next_at = Some(now.saturating_add(i64::from(sink.every_seconds)));
544 }
545 }
546 if manual {
547 host.telemetry
548 .manual
549 .lock()
550 .map_err(|_| Error::Config("telemetry manual request lock poisoned".into()))?
551 .remove(&sink.id);
552 }
553 let host = host.clone();
554 tasks.spawn(async move {
555 let sink = &config.sinks[index];
556 let result = deliver(&host.telemetry, sink, &envelope, 0).await;
557 let http = match &result {
558 Ok((status, _)) => Some(*status),
559 Err(error) => error.status,
560 };
561 let acknowledge = match &result {
562 Ok(_) => true,
563 Err(error) => error.permanent,
564 };
565 let mut error = result.err().map(|error| error.message);
566 if acknowledge
567 && let Some(cursor) = cursor
568 && let Err(failure) = host
569 .advance_telemetry(&sink.id, cursor, config.revision)
570 .await
571 {
572 error = Some(failure.to_string());
573 }
574 match host.telemetry.statuses.lock() {
575 Ok(mut statuses) => {
576 if let Some(status) = statuses.get_mut(&sink.id) {
577 status.in_flight = false;
578 status.last_status = http;
579 if error.is_none() {
580 status.last_success_at = Some(chrono::Utc::now().timestamp());
581 status.consecutive_failures = 0;
582 } else {
583 status.consecutive_failures =
584 status.consecutive_failures.saturating_add(1);
585 }
586 status.last_error = error;
587 }
588 }
589 Err(_) => eprintln!("telemetry delivery status lock poisoned"),
590 }
591 host.telemetry.notify.notify_one();
592 });
593 Ok(())
594 }
595 pub(crate) async fn pending(host: &GatewayHost) -> bool {
596 match host.telemetry_pending().await {
597 Ok(pending) => pending,
598 Err(error) => {
599 eprintln!("telemetry pending check failed: {error}");
600 false
601 }
602 }
603 }
604}
605
606#[derive(Debug, thiserror::Error)]
607#[error("{message}")]
608struct DeliveryError {
609 status: Option<u16>,
610 message: String,
611 permanent: bool,
612}
613impl DeliveryError {
614 fn caused(context: &str, cause: &dyn std::error::Error) -> Self {
615 use std::fmt::Write as _;
616 let mut message = format!("{context}: {cause}");
617 let mut source = cause.source();
618 while let Some(cause) = source {
619 let _ = write!(message, ": {cause}");
621 source = cause.source();
622 }
623 Self {
624 status: None,
625 message,
626 permanent: false,
627 }
628 }
629 fn transient(message: &str) -> Self {
630 Self {
631 status: None,
632 message: message.into(),
633 permanent: false,
634 }
635 }
636}
637async fn deliver(
638 telemetry: &Telemetry,
639 sink: &TelemetrySink,
640 envelope: &Value,
641 response_limit: usize,
642) -> std::result::Result<(u16, Vec<u8>), DeliveryError> {
643 let error = DeliveryError::transient;
644 let client = telemetry.client.as_ref().map_err(|cause| {
645 error(&format!(
646 "telemetry HTTP client initialization failed: {cause}"
647 ))
648 })?;
649 let mut request = match sink.method {
650 SinkMethod::Post => {
651 let bytes = serde_json::to_vec(envelope)
652 .map_err(|cause| error(&format!("telemetry encoding failed: {cause}")))?;
653 if bytes.len() > 64 * 1024 {
654 return Err(DeliveryError {
655 status: None,
656 message: "telemetry envelope exceeds 64 KiB".into(),
657 permanent: true,
658 });
659 }
660 client
661 .post(&sink.url)
662 .header("content-type", "application/json")
663 .body(bytes)
664 }
665 SinkMethod::Get => client.get(&sink.url),
666 };
667 for (name, value) in &sink.headers {
668 request = request.header(name, value);
669 }
670 let token = if let Some(name) = &sink.bearer_env {
671 Some(std::env::var(name).map_err(|_| error("bearer environment variable unavailable"))?)
672 } else if let Some(path) = &sink.bearer_file {
673 let path = telemetry.state_dir.join(path);
674 let state_dir = &telemetry.state_dir;
675 let canonical = tokio::fs::canonicalize(&path)
676 .await
677 .map_err(|cause| error(&format!("bearer file unavailable: {cause}")))?;
678 if !canonical.starts_with(state_dir) {
679 return Err(error("bearer file escapes state directory"));
680 }
681 Some(
682 tokio::task::spawn_blocking(move || crate::config::load_secret_file(&path))
683 .await
684 .map_err(|cause| error(&format!("bearer file task failed: {cause}")))?
685 .map_err(|cause| error(&format!("invalid bearer file: {cause}")))?,
686 )
687 } else {
688 None
689 };
690 if let Some(token) = token {
691 if token.is_empty()
692 || token.len() > 16 * 1024
693 || !token.bytes().all(|b| (33..=126).contains(&b))
694 {
695 return Err(error("invalid bearer token"));
696 }
697 request = request.bearer_auth(token);
698 }
699 let mut response = request
701 .send()
702 .await
703 .map_err(|cause| DeliveryError::caused("telemetry request failed", &cause.without_url()))?;
704 let status = response.status();
705 if status.is_success() {
706 let mut body = Vec::new();
707 if response_limit > 0 {
708 while let Some(chunk) = response.chunk().await.map_err(|cause| {
709 DeliveryError::caused("telemetry response failed", &cause.without_url())
710 })? {
711 if chunk.len() > response_limit.saturating_sub(body.len()) {
712 return Err(error("telemetry response exceeds its size limit"));
713 }
714 body.extend_from_slice(&chunk);
715 }
716 }
717 return Ok((status.as_u16(), body));
718 }
719 Err(DeliveryError {
720 status: Some(status.as_u16()),
721 message: format!("telemetry collector returned HTTP {}", status.as_u16()),
722 permanent: status.is_client_error() && status.as_u16() != 408 && status.as_u16() != 429,
723 })
724}