meerkat-rest 0.8.34

REST API server for Meerkat
#![cfg(feature = "integration-real-tests")]
#![allow(unused_mut, clippy::unwrap_used, clippy::expect_used)]
use axum::body::Body;
use axum::http::Request;
use http_body_util::BodyExt;
use meerkat::{
    AgentFactory, Config, FactoryAgentBuilder, MemoryStore, PersistenceBundle,
    PersistentSessionService, SessionStore,
};
use meerkat_client::TestClient;
use meerkat_core::MemoryConfigStore;
#[cfg(feature = "mob")]
use meerkat_mob_mcp::wire_mob_tools;
use meerkat_rest::{AppState, router};
use meerkat_store::StoreAdapter;
use serde_json::json;
use std::sync::Arc;
use tempfile::TempDir;
use tokio::time::{Duration, timeout};
use tower::ServiceExt;

#[tokio::test]
#[ignore = "lane:e2e-live"]
async fn integration_real_live_continue_hangs() {
    let temp_dir = TempDir::new().unwrap();
    let project_root = temp_dir.path().join("project");
    std::fs::create_dir_all(project_root.join(".rkat")).unwrap();

    let config = Config::default();
    let store_path = temp_dir.path().join("sessions");
    let (event_tx, _) = tokio::sync::broadcast::channel(16);

    let store: Arc<dyn SessionStore> = Arc::new(MemoryStore::new());

    let factory = AgentFactory::new(store_path.clone())
        .builtins(true)
        .shell(true)
        .project_root(project_root.clone());
    let provider_registry = factory.provider_runtime_registry();
    let mut builder = FactoryAgentBuilder::new(factory, config.clone());
    builder.default_llm_client = Some(Arc::new(TestClient::default()));
    let persistence = PersistenceBundle::new(
        store,
        Arc::new(meerkat_runtime::InMemoryRuntimeStore::new()),
        Arc::new(meerkat_store::MemoryBlobStore::new()),
    );
    let runtime_adapter = persistence.runtime_adapter();
    let workgraph_store = persistence.workgraph_store();
    builder.default_session_store = Some(Arc::new(StoreAdapter::new(persistence.session_store())));
    #[cfg(feature = "mob")]
    let builder_mob_tools_slot = Arc::clone(&builder.default_mob_tools);
    let (session_store_inner, runtime_store, blob_store) = persistence.into_parts();
    let mut session_service =
        PersistentSessionService::new(builder, 100, session_store_inner, runtime_store, blob_store);
    let session_service = Arc::new(session_service);
    #[cfg(feature = "mob")]
    let mob_state = wire_mob_tools(
        &builder_mob_tools_slot,
        session_service.clone(),
        Some(runtime_adapter.clone()),
        None,
        meerkat_mob::MobControlPrincipal::Owner,
    );
    let config_store: Arc<dyn meerkat_core::ConfigStore> = Arc::new(MemoryConfigStore::new(
        config.clone(),
        meerkat_models::canonical(),
    ));
    let config_runtime = Arc::new(meerkat_core::ConfigRuntime::new(
        Arc::clone(&config_store),
        store_path.join("config_state.json"),
    ));

    let state = AppState {
        store_path: store_path.clone(),
        max_tokens: config.agent.resolved_max_tokens_per_turn(),
        rest_host: config.rest.host.clone().into(),
        rest_port: config.rest.port,
        enable_builtins: true,
        enable_shell: true,
        project_root: Some(project_root.clone()),
        context_root: None,
        user_config_root: None,
        llm_client_override: Some(Arc::new(TestClient::default())),
        config_store,
        event_tx,
        session_service,
        schedule_service: meerkat::ScheduleService::new(Arc::new(
            meerkat::MemoryScheduleStore::default(),
        )),
        workgraph_service: meerkat::WorkGraphService::with_scope(
            workgraph_store,
            "test-realm",
            meerkat::WorkNamespace::default(),
        ),
        webhook_auth: meerkat_rest::webhook::WebhookAuth::None,
        realm: meerkat_core::RealmId::parse("test-realm").expect("valid realm"),
        realm_config_source: Arc::new(meerkat_store::FilesystemRealmConfigSource::new(
            temp_dir.path().join("empty-realms"),
            temp_dir.path().join("empty-realms").join("__no_global__"),
            meerkat_models::canonical(),
        )),
        instance_id: None,
        backend: "sqlite".to_string(),
        resolved_paths: meerkat_core::ConfigResolvedPaths {
            root: store_path.display().to_string(),
            manifest_path: String::new(),
            config_path: String::new(),
            sessions_sqlite_path: None,
            sessions_jsonl_dir: String::new(),
        },
        expose_paths: false,
        config_runtime,
        realm_lease: Arc::new(tokio::sync::Mutex::new(None)),
        skill_runtime: None,
        runtime_adapter: runtime_adapter.clone(),
        runtime_pre_admissions: meerkat_rest::default_rest_runtime_pre_admissions(),
        runtime_registration_locks: meerkat_rest::default_rest_runtime_registration_locks(),
        schedule_host: Arc::default(),
        request_executor: std::sync::Arc::new(meerkat::surface::SurfaceRequestExecutor::new(
            std::time::Duration::from_secs(5),
        )),
        #[cfg(feature = "mob")]
        mob_state,
        #[cfg(feature = "mcp")]
        mcp_sessions: Arc::new(tokio::sync::RwLock::new(std::collections::HashMap::new())),
        provider_auth_persistence: meerkat_providers::auth_store::ProviderAuthPersistence::new(
            Arc::new(meerkat_providers::auth_store::EphemeralTokenStore::new()),
            Arc::new(meerkat_providers::auth_store::InMemoryCoordinator::new()),
        ),
        auth_lease: runtime_adapter.generated_auth_lease_handle(),
        provider_registry,
    };

    let app = router(state);
    let run_payload = json!({"prompt":"hi"});
    let req1 = Request::builder()
        .method("POST")
        .uri("/sessions")
        .header("content-type", "application/json")
        .body(Body::from(serde_json::to_vec(&run_payload).unwrap()))
        .unwrap();
    let resp1 = timeout(Duration::from_secs(5), app.clone().oneshot(req1))
        .await
        .unwrap()
        .unwrap();
    let body1 = resp1.into_body().collect().await.unwrap().to_bytes();
    let run_json: serde_json::Value = serde_json::from_slice(&body1).unwrap();
    let sid = run_json["session_id"].as_str().unwrap().to_string();

    let continue_payload = json!({"session_id": sid, "prompt":"next"});
    let sid2 = continue_payload["session_id"].as_str().unwrap();
    let req2 = Request::builder()
        .method("POST")
        .uri(format!("/sessions/{sid2}/messages"))
        .header("content-type", "application/json")
        .body(Body::from(serde_json::to_vec(&continue_payload).unwrap()))
        .unwrap();
    let res2 = timeout(Duration::from_secs(2), app.oneshot(req2)).await;
    assert!(
        res2.is_ok(),
        "continue timed out: likely hung waiting event forwarder"
    );
}