mobius-gateway 0.16.15

Headless authenticated gateway for möbius frontends
Documentation
use super::*;
use crate::telemetry::{TelemetryConfig, TelemetrySection, TelemetrySink, TelemetrySinkReport};
use serde_json::{Value, json};

#[derive(Default)]
pub(super) struct StorageMeasurements {
    completed: Mutex<Option<CompletedStorageMeasurement>>,
    revision: AtomicU64,
}

struct CompletedStorageMeasurement {
    completed_at: std::time::Instant,
    revision: u64,
    usage: crate::storage_usage::StorageUsage,
}

impl StorageMeasurements {
    pub(super) fn invalidate(&self) {
        self.revision.fetch_add(1, Ordering::AcqRel);
    }
}

impl GatewayHost {
    pub(crate) fn invalidate_storage_usage(&self) {
        self.storage_reads.invalidate();
    }

    pub(crate) async fn configure_telemetry(
        &self,
        expected: u64,
        mut sinks: Vec<TelemetrySink>,
        preserve_auth: &[String],
    ) -> std::result::Result<(), Rejection> {
        let _mutation = self.begin_mutation().await?;
        let state = self.state.lock().await;
        let mut live = state
            .config
            .lock()
            .map_err(|_| internal("configuration lock poisoned"))?;
        if live.telemetry.revision != expected {
            return Err(invalid_config(
                "telemetry revision changed; refresh and retry",
            ));
        }
        if preserve_auth.len() > 16 {
            return Err(invalid_config("too many preserved credentials"));
        }
        for id in preserve_auth {
            let previous = live
                .telemetry
                .sinks
                .iter()
                .find(|sink| &sink.id == id)
                .ok_or_else(|| invalid_config("cannot preserve credentials for an unknown sink"))?;
            let replacement = sinks
                .iter_mut()
                .find(|sink| &sink.id == id)
                .ok_or_else(|| invalid_config("preserved sink is missing"))?;
            if replacement.url != previous.url {
                return Err(invalid_config(
                    "cannot preserve credentials for a changed destination",
                ));
            }
            if replacement.bearer_env.is_some() || replacement.bearer_file.is_some() {
                return Err(invalid_config(
                    "cannot both preserve and replace credentials",
                ));
            }
            replacement.bearer_env.clone_from(&previous.bearer_env);
            replacement.bearer_file.clone_from(&previous.bearer_file);
        }
        let mut config = live.clone();
        config.telemetry = TelemetryConfig {
            policy: live.telemetry.policy,
            revision: expected
                .checked_add(1)
                .ok_or_else(|| invalid_config("telemetry revision overflow"))?,
            sinks,
        };
        config.validate().map_err(invalid_config)?;
        state.store.save(&config).map_err(internal)?;
        if let Err(error) = state.bots.sync_telemetry_cursors(&config.telemetry.sinks) {
            state.store.save(&live).map_err(internal)?;
            return Err(internal(error));
        }
        self.telemetry
            .configure(config.telemetry.clone())
            .map_err(internal)?;
        *live = config;
        Ok(())
    }

    pub(crate) async fn schedule_telemetry(&self, id: String) -> Result<()> {
        self.telemetry.request_manual(id)?;
        self.telemetry.notify.notify_one();
        Ok(())
    }

    pub(crate) async fn telemetry_pending(&self) -> Result<bool> {
        let config = self.telemetry.config()?;
        let state = self.state.lock().await;
        for sink in config
            .sinks
            .iter()
            .filter(|sink| sink.enabled && !sink.events.is_empty())
        {
            if self.telemetry.can_drain(&sink.id)? && state.bots.telemetry_count(sink)? > 0 {
                return Ok(true);
            }
        }
        Ok(false)
    }

    pub(crate) async fn storage_usage_request(
        &self,
    ) -> std::result::Result<crate::storage_usage::StorageUsage, Rejection> {
        // The async guard serializes only the bounded walk, so simultaneous clients
        // share its completed result instead of starting competing disk scans.
        let mut completed = tokio::time::timeout(
            std::time::Duration::from_secs(2),
            self.storage_reads.completed.lock(),
        )
        .await
        .map_err(|_| internal("storage measurement timed out waiting for an active measurement"))?;
        let revision = self.storage_reads.revision.load(Ordering::Acquire);
        if let Some(previous) = completed.as_ref()
            && previous.revision == revision
            && previous.completed_at.elapsed() < std::time::Duration::from_secs(5)
        {
            // Both the retained snapshot and each wire response require owned data.
            return Ok(previous.usage.clone());
        }
        let usage = self.storage_usage().await.map_err(internal)?;
        *completed = Some(CompletedStorageMeasurement {
            completed_at: std::time::Instant::now(),
            revision,
            usage: usage.clone(),
        });
        Ok(usage)
    }

    pub(crate) async fn telemetry_report(&self) -> Result<(u64, Vec<TelemetrySinkReport>)> {
        let config = self.telemetry.config()?;
        let state = self.state.lock().await;
        let mut reports = Vec::new();
        for sink in &config.sinks {
            let mut sink = sink.clone();
            let mut status = self.telemetry.status(&sink.id)?;
            if !sink.events.is_empty() {
                status.events_pending = state.bots.telemetry_count(&sink)?;
            }
            let auth = if sink.bearer_env.is_some() {
                crate::telemetry::SinkAuth::BearerEnv
            } else if sink.bearer_file.is_some() {
                crate::telemetry::SinkAuth::BearerFile
            } else {
                crate::telemetry::SinkAuth::None
            };
            sink.redact_report()?;
            reports.push(TelemetrySinkReport { sink, auth, status });
        }
        Ok((config.revision, reports))
    }

    pub(crate) async fn telemetry_snapshot(
        &self,
        sections: &[TelemetrySection],
        connected_clients: usize,
    ) -> Result<Value> {
        let mut snapshot = json!({"machine_name": local_machine_name()?});
        if sections.is_empty() {
            return Ok(snapshot);
        }
        if sections.contains(&TelemetrySection::Activity) {
            let activity = self
                .runtime_activity()
                .await
                .map_err(|e| Error::Config(e.message))?;
            snapshot["activity"] = json!({"idle": activity.idle && connected_clients == 0, "connected_clients": connected_clients, "activity_revision": activity.activity_revision, "next_routine_at": activity.next_routine_at, "active_sessions": activity.active_sessions, "running_routines": activity.running_routines});
        }
        if sections.contains(&TelemetrySection::Storage) {
            let mut usage = self.storage_usage().await?;
            usage.bound_telemetry_details();
            snapshot["storage"] = serde_json::to_value(usage)?;
        }
        if sections.contains(&TelemetrySection::Usage) {
            let state = self.state.lock().await;
            let config = state
                .config
                .lock()
                .map_err(|_| Error::Config("configuration lock poisoned".into()))?;
            let cutoff = u64::try_from(Utc::now().timestamp().max(0)).unwrap_or_default() / 86_400;
            snapshot["usage"] = serde_json::to_value(
                config
                    .profile()
                    .daily_usage
                    .into_iter()
                    .filter(|day| day.unix_day >= cutoff.saturating_sub(30))
                    .collect::<Vec<_>>(),
            )?;
        }
        if sections.contains(&TelemetrySection::Runs) {
            let checkpoints = Arc::clone(&self.state.lock().await.checkpoints);
            snapshot["runs"] = serde_json::to_value(
                gateway_run_stats(&gateway_session_summaries(&checkpoints).await?)?.completed,
            )?;
        }
        Ok(snapshot)
    }

    pub(crate) async fn telemetry_events(
        &self,
        sink: &TelemetrySink,
    ) -> Result<(Vec<Value>, Option<i64>, u64)> {
        let (batch, pending, checkpoints, activities, bots) = {
            let state = self.state.lock().await;
            let (batch, pending) = state.bots.telemetry_batch(sink)?;
            (
                batch,
                pending,
                Arc::clone(&state.checkpoints),
                Arc::clone(&state.activities),
                Arc::clone(&state.bots),
            )
        };
        let mut cursor = None;
        let mut bytes = 0_usize;
        let sessions = if batch.is_empty() {
            Vec::new()
        } else {
            session_catalog(&checkpoints, &activities).await?
        };
        let mut events = Vec::with_capacity(batch.len());
        for (row, event) in batch {
            let mut value = serde_json::to_value(&event)?;
            if let crate::wire::HookSource::Session { session_id } = &event.source
                && let Some(session) = sessions.iter().find(|s| &s.session_id == session_id)
            {
                let bot = bots.find_bot(&session.session_context.owner_id)?;
                let mut enrichment =
                    json!({"title": session.title, "run_count": session.execution_stats.run_count});
                if let Some(bot) = bot {
                    enrichment["bot_name"] = json!(bot.name);
                }
                if matches!(
                    event.data,
                    crate::wire::HookData::SessionTurnFinished {
                        outcome: ExecutionOutcome::Completed,
                        ..
                    }
                ) {
                    enrichment["final_message"] = json!(
                        session
                            .activity
                            .message
                            .as_deref()
                            .unwrap_or_default()
                            .chars()
                            .take(1000)
                            .collect::<String>()
                    );
                }
                value["session"] = enrichment;
            }
            let size = serde_json::to_vec(&value)?.len();
            if !events.is_empty() && bytes.saturating_add(size) > 48 * 1024 {
                break;
            }
            bytes = bytes.saturating_add(size);
            cursor = Some(row);
            events.push(value);
        }
        Ok((events, cursor, pending))
    }

    pub(crate) async fn advance_telemetry(
        &self,
        id: &str,
        cursor: i64,
        revision: u64,
    ) -> Result<()> {
        let state = self.state.lock().await;
        if state
            .config
            .lock()
            .map_err(|_| Error::Config("configuration lock poisoned".into()))?
            .telemetry
            .revision
            != revision
        {
            return Ok(());
        }
        state.bots.advance_telemetry(id, cursor)
    }
}

impl GatewayHost {
    pub(crate) async fn storage_usage(&self) -> Result<crate::storage_usage::StorageUsage> {
        let budget = crate::storage_usage::MeasurementBudget::new();
        let (root, checkpoints, files, bots, limit_bytes, tls) = {
            let state = tokio::time::timeout_at(budget.deadline(), self.state.lock())
                .await
                .map_err(|_| {
                    Error::Config("storage measurement timed out waiting for gateway state".into())
                })?;
            let config = state
                .config
                .lock()
                .map_err(|_| Error::Config("configuration lock poisoned".into()))?;
            (
                state.store.state_dir().to_path_buf(),
                Arc::clone(&state.checkpoints),
                state.session_files.clone(),
                Arc::clone(&state.bots),
                config.runtime.storage_limit_bytes,
                config.tls.clone(),
            )
        };
        crate::storage_usage::measure_usage(
            root,
            checkpoints,
            files,
            bots,
            limit_bytes,
            tls,
            budget,
        )
        .await
    }
}