use std::net::{IpAddr, Ipv4Addr, SocketAddr}; use axum::{ body::{Body, to_bytes}, http::{Request, StatusCode, header}, }; use crank_observability::{ DURATION_BUCKETS_SECONDS, MetricsConfig, MetricsConfigError, MetricsSurface, ServiceIdentity, metric_schema, }; use tokio::io::{AsyncReadExt, AsyncWriteExt}; use tower::ServiceExt; fn identity() -> ServiceIdentity { ServiceIdentity::try_new("admin-api", "0.3.1", "test").expect("valid identity") } #[test] fn loopback_is_allowed_without_a_token() { let config = MetricsConfig::new( true, SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), 9464), None, ) .expect("loopback metrics must be safe by default"); assert_eq!(config.bind_addr().to_string(), "127.0.0.1:9464"); assert!(!config.requires_authentication()); } #[test] fn disabled_config_remains_disabled_without_auth_side_effects() { let config = MetricsConfig::new( false, SocketAddr::new(IpAddr::V4(Ipv4Addr::UNSPECIFIED), 9464), None, ) .expect("disabled listener does not require a token"); assert!(!config.enabled()); } #[tokio::test] async fn bind_collision_returns_a_typed_safe_error() { let listener = tokio::net::TcpListener::bind((Ipv4Addr::LOCALHOST, 0)) .await .unwrap(); let address = listener.local_addr().unwrap(); let config = MetricsConfig::new(true, address, None).unwrap(); let error = match MetricsSurface::for_test(config, identity()) .unwrap() .bind() .await { Ok(_) => panic!("occupied port must fail closed"), Err(error) => error, }; assert_eq!(error.to_string(), "failed to bind metrics listener"); } #[tokio::test] async fn default_service_ports_bind_and_serve_real_metrics_listeners() { for (service, port) in [("admin-api", 9464), ("mcp-server", 9465)] { let config = MetricsConfig::new( true, SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), port), None, ) .expect("default loopback metrics config"); let identity = ServiceIdentity::try_new(service, "0.3.1", "test").expect("valid service identity"); let server = MetricsSurface::for_test(config, identity) .expect("metrics surface") .bind() .await .expect("default metrics port must bind"); assert_eq!(server.local_addr().expect("local address").port(), port); let task = tokio::spawn(server.serve()); let response = raw_get(port, "/metrics").await; assert!(response.starts_with("HTTP/1.1 200 OK"), "{response}"); assert!(response.contains("content-type: text/plain")); task.abort(); let _ = task.await; } } #[test] fn non_loopback_without_a_token_is_rejected_without_secret_data() { let error = MetricsConfig::new( true, SocketAddr::new(IpAddr::V4(Ipv4Addr::UNSPECIFIED), 9464), None, ) .expect_err("external metrics must require authentication"); assert!(matches!( error, MetricsConfigError::MissingTokenForExternalBind )); assert!(!error.to_string().contains("token=")); } #[test] fn metrics_surface_rejects_an_unregistered_process_identity() { let config = MetricsConfig::new( true, SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), 9464), None, ) .unwrap(); let identity = ServiceIdentity::try_new("user-controlled", "0.3.1", "test").unwrap(); assert!(MetricsSurface::for_test(config, identity).is_err()); } #[test] fn schema_is_closed_and_uses_fixed_duration_buckets() { let schema = metric_schema(); let names: Vec<_> = schema.iter().map(|metric| metric.name).collect(); assert!(names.contains(&"crank_http_requests_total")); assert!(names.contains(&"crank_http_request_duration_seconds")); assert!(names.contains(&"crank_mcp_requests_total")); assert!(names.contains(&"crank_mcp_active_sessions")); assert!(names.contains(&"crank_mcp_active_streams")); assert!(names.contains(&"crank_tool_invocations_total")); assert!(names.contains(&"crank_runtime_inflight")); assert!(names.contains(&"crank_runtime_cache_total")); assert!(names.contains(&"crank_idempotency_total")); assert!(names.contains(&"crank_confirmation_total")); assert!(names.contains(&"crank_db_pool_connections")); assert!(names.contains(&"crank_catalog_tools")); assert!(names.contains(&"crank_invocation_history_lost_total")); assert!(names.contains(&"crank_telemetry_export_failures_total")); for metric in schema { for forbidden in [ "workspace", "agent_id", "operation_id", "request_id", "trace_id", "correlation_id", "url", "error_message", "text", ] { assert!( !metric.labels.contains(&forbidden), "{} exposes forbidden label {forbidden}", metric.name ); } } assert_eq!( DURATION_BUCKETS_SECONDS, &[ 0.005, 0.01, 0.025, 0.05, 0.1, 0.25, 0.5, 1.0, 2.5, 5.0, 10.0, 30.0, 60.0 ] ); } #[tokio::test] async fn external_surface_protects_both_routes_and_exposes_nothing_else() { let token = "metrics-canary-secret"; let config = MetricsConfig::new( true, SocketAddr::new(IpAddr::V4(Ipv4Addr::UNSPECIFIED), 9464), Some(token.to_owned()), ) .expect("external metrics with token"); let surface = MetricsSurface::for_test(config, identity()).expect("test metrics surface"); let app = surface.router(); for path in ["/metrics", "/health"] { let unauthorized = app .clone() .oneshot(Request::get(path).body(Body::empty()).expect("request")) .await .expect("response"); assert_eq!(unauthorized.status(), StatusCode::UNAUTHORIZED); let authorized = app .clone() .oneshot( Request::get(path) .header(header::AUTHORIZATION, format!("bEaReR {token}")) .body(Body::empty()) .expect("request"), ) .await .expect("response"); assert_eq!(authorized.status(), StatusCode::OK); let body = to_bytes(authorized.into_body(), 1024 * 1024) .await .expect("bounded body"); assert!(!String::from_utf8_lossy(&body).contains(token)); } for authorization in [None, Some(format!("Bearer {token}"))] { let mut request = Request::get("/api/operations"); if let Some(authorization) = authorization { request = request.header(header::AUTHORIZATION, authorization); } let absent = app .clone() .oneshot(request.body(Body::empty()).expect("request")) .await .expect("response"); assert_eq!(absent.status(), StatusCode::NOT_FOUND); } } async fn raw_get(port: u16, path: &str) -> String { let mut stream = tokio::net::TcpStream::connect((Ipv4Addr::LOCALHOST, port)) .await .expect("metrics listener connection"); stream .write_all( format!("GET {path} HTTP/1.1\r\nHost: 127.0.0.1:{port}\r\nConnection: close\r\n\r\n") .as_bytes(), ) .await .expect("metrics request"); let mut response = Vec::new(); stream .read_to_end(&mut response) .await .expect("metrics response"); String::from_utf8(response).expect("utf-8 response") }