mobius-gateway 0.16.25

Headless authenticated gateway for möbius frontends
Documentation
use mobius::backend::checkpoint::Checkpoint;
use mobius::protocol::{SessionContext, TokenUsage};

use super::*;
use mobius::backend::checkpoint::ExecutionStats;

mod activity;
mod bots;
mod descriptors;
mod hooks;
mod lifecycle;
mod live_chats;
mod projection;
mod replay;

impl GatewayHost {
    pub(crate) async fn pause_activity_reply_for_test(
        &self,
    ) -> (oneshot::Receiver<()>, oneshot::Sender<()>, JoinHandle<()>) {
        let (commands, mut receiver) = mpsc::channel(2);
        let (events, _) = broadcast::channel(1);
        let alive = Arc::new(AtomicBool::new(true));
        let terminated = Arc::new(AtomicBool::new(false));
        let termination = Arc::new(tokio::sync::Notify::new());
        self.state.lock().await.sessions.insert(
            "activity-probe".into(),
            HostHandle {
                inner: Arc::new(HostInner {
                    session_id: "activity-probe".into(),
                    commands,
                    events,
                    alive: Arc::clone(&alive),
                    terminated: Arc::clone(&terminated),
                    termination: Arc::clone(&termination),
                    session_mutations: Arc::new(RwLock::new(())),
                    realtime_voice: Arc::new(Mutex::new(())),
                    gateway_sandbox: std::sync::Weak::new(),
                }),
            },
        );
        let (entered, waiting) = oneshot::channel();
        let (release, mut released) = oneshot::channel();
        let task = tokio::spawn(async move {
            let mut entered = Some(entered);
            let mut held_reply = None;
            loop {
                tokio::select! {
                    command = receiver.recv() => match command {
                        Some(HostCommand::RuntimeIsIdle { reply }) => {
                            if let Some(entered) = entered.take() {
                                held_reply = Some(reply);
                                let _ = entered.send(());
                            } else {
                                let _ = reply.send(Ok(false));
                            }
                        }
                        Some(HostCommand::Shutdown) | None => break,
                        _ => panic!("unexpected activity probe command"),
                    },
                    _ = &mut released => break,
                }
            }
            drop(held_reply);
            alive.store(false, Ordering::Release);
            terminated.store(true, Ordering::Release);
            termination.notify_waiters();
        });
        (waiting, release, task)
    }
}

pub(crate) async fn ensure_test_bot(
    gateway: &GatewayHost,
) -> std::result::Result<crate::wire::BotRecord, Rejection> {
    let state = gateway.state.lock().await;
    if let Some(bot) = state.bots.bots().map_err(internal)?.into_iter().next() {
        return Ok(bot);
    }
    let mut config = state.config()?;
    if config.bot_defaults.is_none() {
        let next = config
            .registering_provider(
                AgentComposition::default().provider,
                "Test".into(),
                Default::default(),
                Vec::new(),
                Vec::new(),
                Vec::new(),
            )
            .map_err(invalid_config)?;
        state.store.save(&next).map_err(internal)?;
        *config = next;
    }
    let composition = config
        .bot_defaults
        .as_ref()
        .expect("provider registration installs Bot defaults")
        .config
        .clone();
    drop(config);
    state
        .bots
        .create_bot("Test Bot", "Own gateway test work.", composition)
        .map_err(invalid_bot)
}

pub(crate) async fn create_test_session(
    gateway: &GatewayHost,
    workspace: &Path,
) -> std::result::Result<HostHandle, Rejection> {
    let bot = ensure_test_bot(gateway).await?;
    gateway.create_session(workspace, &bot.id).await
}

pub(crate) fn timer_definition(
    workspace: &Path,
    instructions: &str,
    schedule: crate::wire::RoutineSchedule,
    ends_at: Option<i64>,
) -> crate::wire::RoutineDefinition {
    crate::wire::RoutineDefinition {
        workspace: workspace.into(),
        instructions: instructions.into(),
        bindings: vec![crate::wire::RoutineBinding {
            id: Uuid::new_v4().to_string(),
            on: crate::wire::HookSelector::Schedule { schedule, ends_at },
            action: crate::wire::RoutineAction::Start,
        }],
    }
}
async fn start_routine(gateway: &GatewayHost, id: &str) -> std::result::Result<(), Rejection> {
    gateway
        .execute_routine_command(
            &crate::wire::RoutineCommand {
                routine_id: id.into(),
                action: crate::wire::RoutineAction::Start,
            },
            None,
            None,
            &Uuid::new_v4().to_string(),
        )
        .await
}