#[path = "integration/common.rs"] mod common; use std::{collections::BTreeSet, time::Duration}; use common::*; use crank_core::PlatformApiKeyScope; use crank_registry::PublishRequest; use metrics_util::debugging::{DebugValue, DebuggingRecorder, Snapshotter}; use serde_json::json; use time::{OffsetDateTime, format_description::well_known::Rfc3339}; #[tokio::test] async fn real_mcp_runtime_upstream_postgres_and_catalog_paths_emit_bounded_metrics() { let recorder = DebuggingRecorder::new(); let snapshotter = recorder.snapshotter(); recorder .install() .expect("isolated integration test recorder"); let registry = test_registry().await; let pool_sampler = mcp_server::pool_metrics::spawn_postgres_pool_metrics(registry.pool().clone()); let upstream_base_url = spawn_upstream_server().await; let operation = test_operation(&upstream_base_url, "crm_create_lead"); registry .create_operation(&test_workspace_id(), &operation, Some("metrics-test")) .await .unwrap(); registry .publish_operation(PublishRequest { workspace_id: &test_workspace_id(), operation_id: &operation.id, version: 1, published_at: &OffsetDateTime::parse("2026-07-31T00:00:00Z", &Rfc3339).unwrap(), published_by: Some("metrics-test"), }) .await .unwrap(); publish_agent_for_operation(®istry, &operation, "metrics-agent").await; let api_key = create_platform_api_key( ®istry, "metrics-agent", "metrics-client", &[PlatformApiKeyScope::Read, PlatformApiKeyScope::Write], ) .await; let mut varied_products = Vec::new(); for index in 0..3 { let upstream_url = spawn_upstream_server().await; let operation_name = format!("customer_operation_{index}"); let operation = test_operation(&upstream_url, &operation_name); registry .create_operation(&test_workspace_id(), &operation, Some("metrics-test")) .await .unwrap(); registry .publish_operation(PublishRequest { workspace_id: &test_workspace_id(), operation_id: &operation.id, version: 1, published_at: &OffsetDateTime::parse("2026-07-31T00:00:00Z", &Rfc3339).unwrap(), published_by: Some("metrics-test"), }) .await .unwrap(); let agent_slug = format!("customer-agent-{index}"); publish_agent_for_operation(®istry, &operation, &agent_slug).await; let key = create_platform_api_key( ®istry, &agent_slug, &format!("metrics-client-{index}"), &[PlatformApiKeyScope::Read, PlatformApiKeyScope::Write], ) .await; varied_products.push((agent_slug, key, operation_name, upstream_url)); } let base_url = spawn_mcp_server(build_test_app( registry.clone(), Duration::ZERO, Some("https://crank.example.com".to_owned()), )) .await; let client = reqwest::Client::new(); let mcp_url = agent_mcp_url(&base_url, "metrics-agent"); let malformed = client .post(&mcp_url) .header(reqwest::header::ACCEPT, "application/json") .header(reqwest::header::AUTHORIZATION, format!("Bearer {api_key}")) .header(reqwest::header::CONTENT_TYPE, "application/json") .body("{") .send() .await .unwrap(); assert_eq!(malformed.status(), reqwest::StatusCode::BAD_REQUEST); let invalid_session_header = client .post(&mcp_url) .header( reqwest::header::ACCEPT, "application/json, text/event-stream", ) .header(reqwest::header::AUTHORIZATION, format!("Bearer {api_key}")) .header( "MCP-Session-Id", reqwest::header::HeaderValue::from_bytes(b"\xff").unwrap(), ) .json(&json!({ "jsonrpc": "2.0", "id": 1, "method": "initialize", "params": {"protocolVersion": "2025-11-25", "capabilities": {}} })) .send() .await .unwrap(); assert_eq!( invalid_session_header.status(), reqwest::StatusCode::BAD_REQUEST ); let session_id = initialize_session(&client, &mcp_url, &api_key).await; assert_gauge(&snapshotter, "crank_mcp_active_sessions", 1.0).await; let _ = post_jsonrpc( &client, &mcp_url, &api_key, Some(&session_id), json!({"jsonrpc": "2.0", "id": 2, "method": "tools/list", "params": {}}), ) .await; call_tool( &client, &mcp_url, &api_key, &session_id, "crm_create_lead", "first@example.com", "req-first", ) .await; delete_session(&client, &mcp_url, &api_key, &session_id).await; assert_gauge(&snapshotter, "crank_mcp_active_sessions", 0.0).await; let session_id = initialize_session(&client, &mcp_url, &api_key).await; assert_gauge(&snapshotter, "crank_mcp_active_sessions", 1.0).await; for (index, (agent_slug, key, operation_name, _)) in varied_products.iter().enumerate() { let url = agent_mcp_url(&base_url, agent_slug); let session = initialize_session(&client, &url, key).await; call_tool( &client, &url, key, &session, operation_name, &format!("variant-{index}@example.com"), &format!("req-product-{index}"), ) .await; delete_session(&client, &url, key, &session).await; } send_invalid_workspace_request(&client, &base_url, &api_key, "customer-workspace-seed").await; tokio::time::sleep(Duration::from_millis(25)).await; let first_snapshot = snapshotter.snapshot().into_vec(); let first_series = series(&first_snapshot); assert!(first_snapshot.iter().any(|(key, _, _, _)| { key.key().name() == "crank_mcp_requests_total" && has_label(key.key(), "method", "invalid") && has_label(key.key(), "outcome", "client_error") })); let names = first_snapshot .iter() .map(|(key, _, _, _)| key.key().name()) .collect::>(); for expected in [ "crank_mcp_requests_total", "crank_mcp_active_sessions", "crank_tool_invocations_total", "crank_tool_invocation_duration_seconds", "crank_upstream_requests_total", "crank_upstream_request_duration_seconds", "crank_db_pool_connections", "crank_catalog_tools", "crank_catalog_estimated_context_tokens", "crank_catalog_warnings", ] { assert!( names.contains(expected), "missing product metric {expected}" ); } for index in 0..32 { call_tool( &client, &mcp_url, &api_key, &session_id, "crm_create_lead", &format!("user-{index}@example.com"), &format!("req-metrics-{index}"), ) .await; send_invalid_workspace_request( &client, &base_url, &api_key, &format!("customer-workspace-{index}"), ) .await; } tokio::time::sleep(Duration::from_millis(25)).await; let diverse_snapshot = snapshotter.snapshot().into_vec(); assert_eq!(series(&diverse_snapshot), first_series); delete_session(&client, &mcp_url, &api_key, &session_id).await; assert_gauge(&snapshotter, "crank_mcp_active_sessions", 0.0).await; let rendered = format!("{diverse_snapshot:?}"); let mut forbidden = vec![ "metrics-agent", "crm_create_lead", "user-31@example.com", "req-metrics-31", &upstream_base_url, ]; for (agent_slug, _, operation_name, upstream_url) in &varied_products { forbidden.extend([ agent_slug.as_str(), operation_name.as_str(), upstream_url.as_str(), ]); } for forbidden in forbidden { assert!( !rendered.contains(forbidden), "untrusted value leaked into metric labels: {forbidden}" ); } pool_sampler.abort(); } async fn delete_session(client: &reqwest::Client, mcp_url: &str, api_key: &str, session_id: &str) { let deleted = client .delete(mcp_url) .header(reqwest::header::AUTHORIZATION, format!("Bearer {api_key}")) .header("MCP-Session-Id", session_id) .send() .await .unwrap(); assert_eq!(deleted.status(), reqwest::StatusCode::NO_CONTENT); } async fn call_tool( client: &reqwest::Client, mcp_url: &str, api_key: &str, session_id: &str, tool_name: &str, email: &str, request_id: &str, ) { let response = post_jsonrpc_response( client, mcp_url, api_key, Some(session_id), Some(request_id), json!({ "jsonrpc": "2.0", "id": request_id, "method": "tools/call", "params": { "name": tool_name, "arguments": {"email": email} } }), ) .await; assert_eq!(response.status(), reqwest::StatusCode::OK); } async fn send_invalid_workspace_request( client: &reqwest::Client, base_url: &str, api_key: &str, workspace_slug: &str, ) { let response = client .post(format!("{base_url}/v1/{workspace_slug}/metrics-agent")) .header(reqwest::header::ACCEPT, "application/json") .header(reqwest::header::AUTHORIZATION, format!("Bearer {api_key}")) .json(&json!({ "jsonrpc": "2.0", "id": 1, "method": "initialize", "params": {"protocolVersion": "2025-11-25", "capabilities": {}} })) .send() .await .unwrap(); assert!(response.status().is_client_error()); } fn series( snapshot: &[( metrics_util::CompositeKey, Option, Option, metrics_util::debugging::DebugValue, )], ) -> BTreeSet { snapshot .iter() .map(|(key, _, _, _)| format!("{key:?}")) .collect() } fn has_label(key: &metrics::Key, name: &str, value: &str) -> bool { key.labels() .any(|label| label.key() == name && label.value() == value) } async fn assert_gauge(snapshotter: &Snapshotter, name: &str, expected: f64) { for _ in 0..50 { let value = snapshotter .snapshot() .into_vec() .into_iter() .find_map(|(key, _, _, value)| { if key.key().name() != name { return None; } match value { DebugValue::Gauge(value) => Some(value.into_inner()), _ => None, } }); if value == Some(expected) { return; } tokio::time::sleep(Duration::from_millis(10)).await; } panic!("{name} did not become {expected}"); }