mobius-gateway 0.16.15

Headless authenticated gateway for möbius frontends
Documentation
use super::args::TelemetryCommand;
use super::*;
use crate::telemetry::{SinkMethod, TelemetrySink};

pub(super) async fn run(
    state_dir: PathBuf,
    command: TelemetryCommand,
    load: fn(&Endpoint) -> Result<Option<String>>,
) -> Result<()> {
    let (store, mut config) = ConfigStore::open(state_dir)?;
    let endpoint = direct_loopback_endpoint(&config)?;
    let running = running_process_record(&store.state_dir().join(PROCESS_FILE))?.is_some();
    if !running {
        ensure_gateway_stopped(&store, &config)?;
    }
    if matches!(command, TelemetryCommand::List) && running {
        let token = load(&endpoint)?
            .ok_or_else(|| Error::Config("local control credential unavailable".into()))?;
        let response = request(&endpoint, &token, |request_id| {
            ClientMessage::GetTelemetry { request_id }
        })
        .await?;
        println!("{}", serde_json::to_string_pretty(&response)?);
        return Ok(());
    }
    let expected_revision = config.telemetry.revision;
    let previous = config.clone();
    match command {
        TelemetryCommand::List => {
            for sink in &mut config.telemetry.sinks {
                sink.redact_report()?;
            }
            println!("{}", serde_json::to_string_pretty(&config.telemetry)?);
            return Ok(());
        }
        TelemetryCommand::Remove { id } => config.telemetry.sinks.retain(|sink| sink.id != id),
        TelemetryCommand::Add {
            id,
            url,
            every_seconds,
            upload_admission,
            sections,
            events,
            bearer_env,
            bearer_file,
            fields,
        } => {
            let fields = fields
                .into_iter()
                .map(|field| {
                    field
                        .split_once('=')
                        .map(|(key, value)| (key.to_owned(), value.to_owned()))
                        .ok_or_else(|| Error::Config("--field requires key=value".into()))
                })
                .collect::<Result<_>>()?;
            let sections = sections
                .into_iter()
                .map(|value| {
                    serde_json::from_value(serde_json::Value::String(value)).map_err(Error::from)
                })
                .collect::<Result<_>>()?;
            let events = events
                .into_iter()
                .map(|value| {
                    serde_json::from_value(serde_json::Value::String(value)).map_err(Error::from)
                })
                .collect::<Result<_>>()?;
            config.telemetry.sinks.retain(|sink| sink.id != id);
            config.telemetry.sinks.push(TelemetrySink {
                id,
                url,
                method: SinkMethod::Post,
                every_seconds,
                upload_admission,
                sections,
                events,
                bearer_env,
                bearer_file,
                fields,
                headers: Default::default(),
                enabled: true,
            });
        }
    }
    config.telemetry.revision = expected_revision
        .checked_add(1)
        .ok_or_else(|| Error::Config("telemetry revision overflow".into()))?;
    config.validate()?;
    if running {
        let token = load(&endpoint)?
            .ok_or_else(|| Error::Config("local control credential unavailable".into()))?;
        request(&endpoint, &token, |request_id| {
            ClientMessage::ConfigureTelemetry {
                request_id,
                expected_revision,
                sinks: config.telemetry.sinks,
                preserve_auth: Vec::new(),
            }
        })
        .await?;
    } else {
        let bots = crate::bots::BotStore::open(store.state_dir())?;
        store.save(&config)?;
        if let Err(error) = bots.sync_telemetry_cursors(&config.telemetry.sinks) {
            store.save(&previous)?;
            return Err(error);
        }
    }
    Ok(())
}

async fn request(
    endpoint: &Endpoint,
    token: &str,
    message: impl FnOnce(String) -> ClientMessage,
) -> Result<ServerMessage> {
    tokio::time::timeout(std::time::Duration::from_secs(15), async {
        let mut options = crate::client::ConnectOptions {
            catalog: crate::wire::CatalogHint {
                skip: crate::wire::READY_SECTIONS.into_iter().collect(),
                ..Default::default()
            },
            ..Default::default()
        };
        let client = GatewayClient::connect_with(endpoint, token, ClientKind::GatewayDashboard, &mut options).await?;
        let (sender, mut events) = client.into_parts();
        let request_id = Uuid::new_v4().to_string();
        sender.send(message(request_id.clone())).await?;
        for _ in 0..MAX_PENDING_FRAMES {
            let frame = events.next().await?.ok_or_else(|| Error::Protocol("gateway disconnected".into()))?;
            if let Some(error) = frame.message.response_error(Some(&request_id)) { return Err(Error::Protocol(error.message.into())); }
            if matches!(&frame.message, ServerMessage::Telemetry { request_id: id, .. } if id == &request_id) { return Ok(frame.message); }
        }
        Err(Error::Protocol("too many unrelated responses".into()))
    }).await.map_err(|_| Error::Protocol("telemetry configuration timed out".into()))?
}