use std::str::FromStr; use std::sync::Arc; use cm_api::AppState; use cm_auth::AuthService; use cm_domain::{Role, SessionKey, User, UserId, Workspace, WorkspaceId}; use cm_llm::ScriptedProvider; use cm_runtime::{Runtime, RuntimeConfig}; use eventsource_stream::Eventsource; use futures::StreamExt; use serde_json::{json, Value}; const SCENARIOS: &str = r#" [[scenario]] marker = "[[scenario:tool-time]]" [[scenario.turns]] events = [ { type = "text", text = "Let me check the clock." }, { type = "tool_use", name = "clock.now", input = {} }, ] [[scenario.turns]] events = [ { type = "text", text = " All done." }, ] "#; struct TestServer { base: String, client: reqwest::Client, } async fn serve(pool: sqlx::PgPool) -> TestServer { let runtime = Runtime::new( pool.clone(), Arc::new(ScriptedProvider::from_toml(SCENARIOS).unwrap()), RuntimeConfig::basic("scripted", 1024), ); let app = cm_api::router(AppState::new(pool, runtime)); let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); let addr = listener.local_addr().unwrap(); tokio::spawn(async move { axum::serve(listener, app).await.unwrap(); }); TestServer { base: format!("http://{addr}"), client: reqwest::Client::new(), } } async fn seed_and_login(pool: &sqlx::PgPool, server: &TestServer) -> (String, String) { let ws = Workspace { id: WorkspaceId::new(), name: "Acme".into(), plan: "team".into(), }; cm_db::repo::workspaces::insert(pool, &ws).await.unwrap(); let owner = User { id: UserId::new(), workspace_id: ws.id, email: format!("{}@acme.test", UserId::new()), role: Role::Owner, display_name: "Owner".into(), created_at: time::OffsetDateTime::UNIX_EPOCH, }; cm_db::repo::users::insert(pool, &owner).await.unwrap(); AuthService::new(pool.clone()) .set_password(owner.id, "pw") .await .unwrap(); let token = server .client .post(format!("{}/api/auth/login", server.base)) .json(&json!({"email": owner.email, "password": "pw"})) .send() .await .unwrap() .json::() .await .unwrap()["token"] .as_str() .unwrap() .to_owned(); let claw: Value = server .client .post(format!("{}/api/claws", server.base)) .bearer_auth(&token) .json(&json!({"name": "Scout", "job_title": "Analyst"})) .send() .await .unwrap() .json() .await .unwrap(); (token, claw["id"].as_str().unwrap().to_owned()) } async fn create_session(server: &TestServer, token: &str, claw_id: &str) -> Value { let res = server .client .post(format!("{}/api/sessions", server.base)) .bearer_auth(token) .json(&json!({"clawId": claw_id, "title": "Research"})) .send() .await .unwrap(); assert_eq!(res.status(), 201); res.json().await.unwrap() } /// Collects SSE events (id, event-name, parsed data) until run completion. async fn collect_sse(response: reqwest::Response) -> Vec<(String, String, Value)> { let mut events = Vec::new(); let mut stream = response.bytes_stream().eventsource(); while let Some(event) = stream.next().await { let event = event.unwrap(); let data: Value = serde_json::from_str(&event.data).unwrap(); let name = event.event.clone(); let done = name == "run_completed" || name == "error"; events.push((event.id, name, data)); if done { break; } } events } #[tokio::test] async fn session_create_returns_a_parseable_session_key() { let pool = cm_testkit::test_pool().await; let server = serve(pool.clone()).await; let (token, claw_id) = seed_and_login(&pool, &server).await; let session = create_session(&server, &token, &claw_id).await; assert_eq!(session["title"], "Research"); let key = SessionKey::from_str(session["sessionKey"].as_str().unwrap()).unwrap(); assert_eq!(key.agent_id.to_string(), claw_id); } #[tokio::test] async fn sessions_list_is_scoped_to_the_claw() { let pool = cm_testkit::test_pool().await; let server = serve(pool.clone()).await; let (token, claw_id) = seed_and_login(&pool, &server).await; create_session(&server, &token, &claw_id).await; create_session(&server, &token, &claw_id).await; let sessions: Value = server .client .get(format!("{}/api/sessions?clawId={claw_id}", server.base)) .bearer_auth(&token) .send() .await .unwrap() .json() .await .unwrap(); assert_eq!(sessions.as_array().unwrap().len(), 2); } #[tokio::test] async fn gateway_streams_a_full_run_with_monotonic_ids() { let pool = cm_testkit::test_pool().await; let server = serve(pool.clone()).await; let (token, claw_id) = seed_and_login(&pool, &server).await; let session = create_session(&server, &token, &claw_id).await; let session_key = session["sessionKey"].as_str().unwrap(); let response = server .client .post(format!("{}/api/gateway?clawId={claw_id}", server.base)) .bearer_auth(&token) .json(&json!({"sessionKey": session_key, "message": "hello"})) .send() .await .unwrap(); assert_eq!(response.status(), 200); assert!(response .headers() .get("content-type") .unwrap() .to_str() .unwrap() .starts_with("text/event-stream")); let events = collect_sse(response).await; assert_eq!(events[0].1, "run_started"); assert_eq!(events.last().unwrap().1, "run_completed"); let ids: Vec = events .iter() .map(|(id, _, _)| id.parse().unwrap()) .collect(); let mut sorted = ids.clone(); sorted.sort_unstable(); assert_eq!(ids, sorted, "ids must be monotonic"); let text: String = events .iter() .filter(|(_, name, _)| name == "text_delta") .map(|(_, _, data)| data["delta"].as_str().unwrap()) .collect(); assert_eq!(text, "I received: hello"); } #[tokio::test] async fn gateway_runs_tools_and_history_replays_with_steps() { let pool = cm_testkit::test_pool().await; let server = serve(pool.clone()).await; let (token, claw_id) = seed_and_login(&pool, &server).await; let session = create_session(&server, &token, &claw_id).await; let session_key = session["sessionKey"].as_str().unwrap(); let response = server .client .post(format!("{}/api/gateway?clawId={claw_id}", server.base)) .bearer_auth(&token) .json(&json!({ "sessionKey": session_key, "message": "time? [[scenario:tool-time]]" })) .send() .await .unwrap(); let events = collect_sse(response).await; assert!(events.iter().any(|(_, name, _)| name == "step_started")); assert!(events .iter() .any(|(_, name, data)| { name == "step_finished" && data["status"] == "ok" })); let history: Value = server .client .get(format!( "{}/api/sessions/history?sessionKey={}&tools=true", server.base, urlencoding::encode(session_key) )) .bearer_auth(&token) .send() .await .unwrap() .json() .await .unwrap(); let messages = history.as_array().unwrap(); assert_eq!(messages.len(), 2); assert_eq!(messages[1]["role"], "agent"); assert_eq!(messages[1]["steps"][0]["tool_name"], "clock.now"); } #[tokio::test] async fn resume_replays_the_exact_journal() { let pool = cm_testkit::test_pool().await; let server = serve(pool.clone()).await; let (token, claw_id) = seed_and_login(&pool, &server).await; let session = create_session(&server, &token, &claw_id).await; let session_key = session["sessionKey"].as_str().unwrap(); let live = collect_sse( server .client .post(format!("{}/api/gateway?clawId={claw_id}", server.base)) .bearer_auth(&token) .json(&json!({"sessionKey": session_key, "message": "ping"})) .send() .await .unwrap(), ) .await; // Reconnect with resumeFrom=0: the same events come back from the // journal, byte-for-byte equal data. let replayed = collect_sse( server .client .post(format!("{}/api/gateway?clawId={claw_id}", server.base)) .bearer_auth(&token) .json(&json!({"sessionKey": session_key, "resumeFrom": 0})) .send() .await .unwrap(), ) .await; assert_eq!(live, replayed); // Partial resume skips already-seen events. let tail = collect_sse( server .client .post(format!("{}/api/gateway?clawId={claw_id}", server.base)) .bearer_auth(&token) .json(&json!({"sessionKey": session_key, "resumeFrom": live[1].0.parse::().unwrap()})) .send() .await .unwrap(), ) .await; assert_eq!(tail.len(), live.len() - 2); } #[tokio::test] async fn gateway_rejects_claws_outside_the_workspace() { let pool = cm_testkit::test_pool().await; let server = serve(pool.clone()).await; let (token, claw_id) = seed_and_login(&pool, &server).await; let (other_token, _) = seed_and_login(&pool, &server).await; let session = create_session(&server, &token, &claw_id).await; let response = server .client .post(format!("{}/api/gateway?clawId={claw_id}", server.base)) .bearer_auth(&other_token) .json(&json!({ "sessionKey": session["sessionKey"], "message": "hi" })) .send() .await .unwrap(); assert_eq!(response.status(), 404); }