#![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"
);
}