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    // Reserve runtime mutation before changing durable configuration, so assignment cannot fail afterward.
283    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        // The report owns its diagnostic strings after the short lock is released.
325        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    /// Serve-loop boundary: collector failures must never stop the gateway.
342    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        // Cancelling the worker drops its JoinSet, cancelling its deliveries too.
384        // Pending event cursors remain unacknowledged and are retried after restart.
385        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            // A broken source or sink cannot starve the remaining destinations.
424            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            // Writing to a String is infallible.
627            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    // Transport diagnostics intentionally omit the URL, which may contain a collector secret.
707    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}