Skip to main content

mobius_gateway/
telemetry.rs

1//! Configured outbound telemetry; no endpoint is enabled by default.
2mod 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/// Durable endpoint configuration.
23#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
24#[serde(default, deny_unknown_fields)]
25pub struct TelemetryConfig {
26    /// Transport settings controlled locally by the gateway operator.
27    pub policy: TelemetryPolicy,
28    /// Local operator command held while runtime activity needs an awake host.
29    #[serde(skip_serializing_if = "Option::is_none")]
30    pub activity_hook: Option<ActivityHookConfig>,
31    /// Optimistic concurrency revision.
32    pub revision: u64,
33    /// Explicitly configured destinations.
34    pub sinks: Vec<TelemetrySink>,
35}
36/// HTTP delivery method.
37#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
38#[serde(rename_all = "snake_case")]
39pub enum SinkMethod {
40    /// JSON envelope request.
41    #[default]
42    Post,
43    /// Bodyless heartbeat request.
44    Get,
45}
46/// Optional snapshot sections.
47#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize)]
48#[serde(rename_all = "snake_case")]
49pub enum TelemetrySection {
50    /// Runtime work and clients.
51    Activity,
52    /// Daily token totals.
53    Usage,
54    /// Completed run totals.
55    Runs,
56    /// Read-only storage totals.
57    Storage,
58    /// Host, gateway process, and container resource measurements.
59    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
73// Reuse each section within one tick, never across measurements or configuration changes.
74struct 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                // Each in-flight request owns its payload after this tick's snapshot is dropped.
96                result[section.key()] = value.clone();
97            }
98        }
99        Ok(result)
100    }
101}
102
103/// One collector subscription. Secret values are never persisted here.
104#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
105#[serde(deny_unknown_fields)]
106pub struct TelemetrySink {
107    /// Unique stable identifier.
108    pub id: String,
109    /// Collector URL.
110    pub url: String,
111    /// Delivery method.
112    #[serde(default)]
113    pub method: SinkMethod,
114    /// Snapshot interval in seconds.
115    pub every_seconds: u32,
116    /// Requested snapshots.
117    #[serde(default)]
118    pub sections: Vec<TelemetrySection>,
119    /// Committed facts to deliver.
120    #[serde(default)]
121    pub events: Vec<HookKind>,
122    /// Non-secret HTTP headers.
123    #[serde(default)]
124    pub headers: BTreeMap<String, String>,
125    /// Environment variable containing a bearer token.
126    #[serde(default)]
127    pub bearer_env: Option<String>,
128    /// Owner-only token file relative to the state directory.
129    #[serde(default)]
130    pub bearer_file: Option<String>,
131    /// Static collector labels.
132    #[serde(default)]
133    pub fields: BTreeMap<String, String>,
134    /// Whether delivery is enabled.
135    #[serde(default = "enabled")]
136    pub enabled: bool,
137    /// Ask this POST collector to admit user uploads; generated files bypass it.
138    #[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        // A collector can carry credentials in its path, query or custom headers.
146        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/// In-memory delivery status.
158#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
159pub struct TelemetrySinkStatus {
160    /// Most recent attempt.
161    pub last_attempt_at: Option<i64>,
162    /// Most recent successful delivery.
163    pub last_success_at: Option<i64>,
164    /// Last HTTP status.
165    pub last_status: Option<u16>,
166    /// Bounded diagnostic, excluding credentials.
167    pub last_error: Option<String>,
168    /// Failures since last success.
169    pub consecutive_failures: u32,
170    /// Next snapshot time.
171    pub next_at: Option<i64>,
172    /// Undelivered matching facts.
173    pub events_pending: u64,
174    /// A request is running.
175    pub in_flight: bool,
176}
177/// Safe endpoint report returned to clients.
178#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
179pub struct TelemetrySinkReport {
180    /// Configuration without credential source names.
181    pub sink: TelemetrySink,
182    /// Credential source kind; never its name or value.
183    pub auth: SinkAuth,
184    /// Current delivery state.
185    pub status: TelemetrySinkStatus,
186}
187
188/// Redacted authentication mechanism exposed to paired clients.
189#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
190#[serde(rename_all = "snake_case")]
191pub enum SinkAuth {
192    /// No bearer credential.
193    None,
194    /// Credential read from the gateway environment.
195    BearerEnv,
196    /// Credential read from a protected file.
197    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    // GatewayConfig owns the persisted configuration; this snapshot is refreshed by
222    // configure_telemetry after saving and shared with in-flight deliveries.
223    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            // One connection pool and root store for the gateway's lifetime.
241            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        // The report owns its diagnostic strings after the short lock is released.
318        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    /// Serve-loop boundary: collector failures must never stop the gateway.
335    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        // Cancelling the worker drops its JoinSet, cancelling its deliveries too.
377        // Pending event cursors remain unacknowledged and are retried after restart.
378        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            // A broken source or sink cannot starve the remaining destinations.
417            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            // Writing to a String is infallible.
620            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    // Transport diagnostics intentionally omit the URL, which may contain a collector secret.
700    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}