mobius 0.16.0

A small, modular Rust framework for building coding agents
Documentation
//! SQLite checkpoint storage tests.

use std::sync::mpsc;

use serde_json::json;
use tokio::sync::oneshot;
use tokio::time::timeout;

use super::*;
use crate::backend::checkpoint::ExecutionOutcome;
use crate::backend::checkpoint::StreamMetrics;
use crate::backend::checkpoint::event_turn_page;
use crate::protocol::EventMsg;
use crate::protocol::ModelStepContentPhase;
use crate::protocol::ModelStepOutcome;
use crate::protocol::TurnAbortedEvent;
use crate::protocol::TurnCompleteEvent;
use crate::protocol::TurnStartedEvent;

fn checkpoint(session_id: impl Into<String>) -> Checkpoint {
    let mut checkpoint = Checkpoint::empty(session_id);
    checkpoint.session_context.owner_id = "test-bot".into();
    checkpoint
}

#[tokio::test]
async fn admission_receipts_commit_with_queue_and_survive_consumption() {
    use crate::backend::checkpoint::{QueuedMessage, QueuedMessageBoundary};
    use crate::protocol::{MessageAuthor, MessageDelivery, MessageEvent};

    let directory = tempfile::tempdir().expect("directory");
    let store = SqliteCheckpoint::new(directory.path().join("receipts.sqlite3")).expect("store");
    let mut state = checkpoint("session");
    store.save(&state, &[], None).await.expect("seed");
    state.pending_messages.push(
        QueuedMessage::new(
            "messages",
            "input",
            QueuedMessageBoundary::Turn,
            MessageEvent {
                author: MessageAuthor::User,
                delivery: MessageDelivery::Turn,
                text: "hello".into(),
                attachments: Vec::new(),
                reply: None,
                message_target: None,
            },
        )
        .expect("message"),
    );
    assert!(store.save(&state, &[], None).await.is_err());
    assert!(
        !store
            .message_accepted("session", "input")
            .await
            .expect("absent receipt")
    );
    state.sequence += 1;
    store.save(&state, &[], None).await.expect("admit");
    state.sequence += 1;
    state.pending_messages.clear();
    store.save(&state, &[], None).await.expect("consume");
    assert!(
        store
            .message_accepted("session", "input")
            .await
            .expect("durable receipt")
    );
    assert!(
        !store
            .message_accepted("other", "input")
            .await
            .expect("session scope")
    );
    store
        .delete_sessions(&["session".into()])
        .await
        .expect("delete");
    assert!(
        !store
            .message_accepted("session", "input")
            .await
            .expect("deleted receipt")
    );
}

fn execution(session_id: &str, turn: u64) -> ExecutionRecord {
    let started_at_ms = i64::try_from(turn * 100).expect("execution start");
    ExecutionRecord {
        session_id: session_id.into(),
        submission_id: format!("submission-{turn}"),
        author: crate::protocol::MessageAuthor::User,
        turn_id: format!("turn-{turn}"),
        started_at_ms,
        finished_at_ms: started_at_ms + 25,
        elapsed_ms: 25,
        outcome: ExecutionOutcome::Completed,
        model_calls: 1,
        tool_calls: turn,
        failed_tool_calls: 0,
        usage: crate::protocol::TokenUsage {
            total_tokens: 1,
            ..crate::protocol::TokenUsage::default()
        },
    }
}

#[path = "sqlite_tests/event_journal.rs"]
mod event_journal;
#[path = "sqlite_tests/journals.rs"]
mod journals;
#[path = "sqlite_tests/model_steps.rs"]
mod model_steps;
#[path = "sqlite_tests/sessions.rs"]
mod sessions;
#[path = "sqlite_tests/strict.rs"]
mod strict;