mobius-cli 0.15.29

The terminal client for a möbius gateway
Documentation
use std::io;

use mobius::{Error, Result};
use mobius_gateway::client::{GatewayEvents, GatewaySender};
use mobius_gateway::wire::{ClientMessage, ReadyPayload, ServerFrame, ServerMessage};
use ratatui::Terminal;
use ratatui::backend::CrosstermBackend;
use ratatui::crossterm::event::Event;
use tokio::time::MissedTickBehavior;
use uuid::Uuid;

use super::render::render;
use super::state::BotsState;
use super::{Action, FollowUp};
use crate::frontend::setup;
use crate::frontend::terminal::{INPUT_POLL, MAX_INPUT_BATCH, poll_event};
use crate::gateway_error;

type BotsTerminal = Terminal<CrosstermBackend<io::Stdout>>;

pub(in crate::frontend) async fn run(
    terminal: &mut BotsTerminal,
    sender: &GatewaySender,
    events: &mut GatewayEvents,
    gateway: &mut ReadyPayload,
    preferred_bot_id: Option<&str>,
    protected_bot_id: Option<&str>,
) -> Result<Option<String>> {
    terminal.clear()?;
    let mut state = BotsState::new(gateway, preferred_bot_id, protected_bot_id);
    request_routines(sender, &mut state).await?;
    let mut tick = tokio::time::interval(INPUT_POLL);
    tick.set_missed_tick_behavior(MissedTickBehavior::Skip);
    let mut events = events.scoped();
    let mut events_open = true;
    let mut dirty = true;

    'screen: loop {
        if dirty {
            terminal.draw(|frame| render(frame, &state, gateway))?;
            dirty = false;
        }
        tokio::select! {
            frame = events.next(), if events_open => {
                let frame = match frame {
                    Ok(frame) => frame,
                    Err(error) => break 'screen Err(gateway_error(error)),
                };
                match frame {
                    Some(frame) => {
                        if let Some(error) = frame.message.response_error(None) {
                            break 'screen Err(Error::Stopped(error.message.into()));
                        }
                        let (follow_up, deferred) =
                            handle_frame(frame.message, gateway, &mut state);
                        if let Some(frame) = deferred {
                            events.defer(frame).map_err(gateway_error)?;
                        }
                        request_follow_up(sender, &mut state, follow_up).await?;
                    }
                    None => {
                        events_open = false;
                        state.fail("Gateway disconnected. Press q to close.");
                    }
                }
                state.clamp(gateway);
                dirty = true;
            }
            _ = tick.tick() => {
                for _ in 0..MAX_INPUT_BATCH {
                    let Some(event) = poll_event()? else { break; };
                    dirty = true;
                    let action = match event {
                        Event::Key(key) => state.handle_key(key, gateway),
                        Event::Paste(value) => {
                            state.paste(&value);
                            Action::None
                        }
                        Event::Resize(_, _)
                        | Event::FocusGained
                        | Event::FocusLost
                        | Event::Mouse(_) => Action::None,
                    };
                    match action {
                        Action::None => {}
                        Action::Exit => break 'screen Ok(None),
                        Action::OpenSession(session_id) => {
                            break 'screen Ok(Some(session_id));
                        }
                        Action::Setup { bot_id, mode } => {
                            if let Err(error) = setup::run_bot(
                                terminal,
                                mode,
                                None,
                                sender,
                                events.reborrow(),
                                gateway,
                                &bot_id,
                            ).await {
                                state.fail(error.to_string());
                            }
                            terminal.clear()?;
                            state.clamp(gateway);
                        }
                        Action::Send {
                            request_id,
                            message,
                            label,
                            follow_up,
                        } => match sender.send(*message).await {
                            Ok(()) => state.begin(request_id, label, follow_up),
                            Err(error) => state.fail(error.to_string()),
                        },
                    }
                }
            }
        }
    }
}

pub(super) fn handle_frame(
    message: ServerMessage,
    gateway: &mut ReadyPayload,
    state: &mut BotsState,
) -> (FollowUp, Option<ServerFrame>) {
    let mut follow_up = FollowUp::None;
    let mut deferred = None;
    match message {
        ServerMessage::Ready { payload } => *gateway = payload,
        ServerMessage::GatewayConfigured {
            request_id,
            payload,
        } => {
            *gateway = payload.clone();
            deferred = Some(ServerFrame::new(ServerMessage::GatewayConfigured {
                request_id,
                payload,
            }));
        }
        ServerMessage::Sessions {
            request_id,
            sessions,
        } => {
            gateway.sessions = sessions.clone();
            if request_id.is_some() {
                deferred = Some(ServerFrame::new(ServerMessage::Sessions {
                    request_id,
                    sessions,
                }));
            }
        }
        ServerMessage::Bots { request_id, bots } => {
            gateway.bots = bots.clone();
            if request_id
                .as_ref()
                .is_some_and(|id| pending_matches(state, id))
            {
                follow_up = state.complete().unwrap_or(FollowUp::None);
            } else if request_id.is_some() {
                deferred = Some(ServerFrame::new(ServerMessage::Bots { request_id, bots }));
            }
        }
        ServerMessage::Routines {
            request_id,
            routines,
        } if pending_matches(state, &request_id) => {
            state.routines = routines;
            follow_up = state.complete().unwrap_or(FollowUp::None);
        }
        ServerMessage::RoutineHistory { request_id, runs }
            if pending_matches(state, &request_id) =>
        {
            state.runs = runs;
            follow_up = state.complete().unwrap_or(FollowUp::None);
        }
        ServerMessage::RoutineRunPreview {
            request_id,
            preview,
        } if pending_matches(state, &request_id) => {
            state.preview = Some(preview);
            follow_up = state.complete().unwrap_or(FollowUp::None);
        }
        ServerMessage::Accepted { request_id } if pending_matches(state, &request_id) => {
            follow_up = state.complete().unwrap_or(FollowUp::None);
        }
        ServerMessage::Rejected {
            request_id,
            message,
            ..
        } if pending_matches(state, &request_id) => state.fail(message),
        message => deferred = Some(ServerFrame::new(message)),
    }
    (follow_up, deferred)
}

fn pending_matches(state: &BotsState, request_id: &str) -> bool {
    state
        .pending
        .as_ref()
        .is_some_and(|pending| pending.request_id == request_id)
}

async fn request_routines(sender: &GatewaySender, state: &mut BotsState) -> Result<()> {
    let request_id = Uuid::new_v4().to_string();
    sender
        .send(ClientMessage::ListRoutines {
            request_id: request_id.clone(),
            bot_id: None,
        })
        .await
        .map_err(gateway_error)?;
    state.begin(request_id, "Load routines", FollowUp::None);
    Ok(())
}

async fn request_runs(
    sender: &GatewaySender,
    state: &mut BotsState,
    routine_id: String,
) -> Result<()> {
    let request_id = Uuid::new_v4().to_string();
    sender
        .send(ClientMessage::ListRoutineHistory {
            request_id: request_id.clone(),
            id: Some(routine_id),
        })
        .await
        .map_err(gateway_error)?;
    state.begin(request_id, "Load run history", FollowUp::None);
    Ok(())
}

async fn request_follow_up(
    sender: &GatewaySender,
    state: &mut BotsState,
    follow_up: FollowUp,
) -> Result<()> {
    match follow_up {
        FollowUp::None => Ok(()),
        FollowUp::Routines => request_routines(sender, state).await,
        FollowUp::Runs(routine_id) => request_runs(sender, state, routine_id).await,
    }
}