use axum::{ Router, middleware, routing::{delete, get, post}, }; use crate::{ auth::{require_session, require_workspace_session}, routes::{ access::{ create_invitation, create_platform_api_key, delete_invitation, delete_platform_api_key, list_invitations, list_memberships, list_platform_api_keys, revoke_platform_api_key, }, agents::{ create_agent, delete_agent, get_agent, get_agent_version, list_agents, publish_agent, save_agent_bindings, update_agent, }, auth::{get_session, login, logout}, auth_profiles::{create_auth_profile, get_auth_profile, list_auth_profiles}, observability::{get_agent_usage, get_log, get_operation_usage, get_usage, list_logs}, operations::{ archive_operation, create_operation, create_version, delete_operation, export_operation, generate_draft, get_operation, get_operation_version, list_grpc_services, list_operations, publish_operation, run_test, update_operation, upload_descriptor_set, upload_input_json, upload_output_json, upload_proto_descriptor, }, workspaces::{create_workspace, get_workspace, list_workspaces, update_workspace}, }, state::AppState, }; pub fn build_app(state: AppState) -> Router { let workspace_router = Router::new() .route("/operations", get(list_operations).post(create_operation)) .route("/operations/import", post(import_operation)) .route( "/operations/{operation_id}", get(get_operation) .patch(update_operation) .delete(delete_operation), ) .route("/operations/{operation_id}/versions", post(create_version)) .route( "/operations/{operation_id}/versions/{version}", get(get_operation_version), ) .route( "/operations/{operation_id}/publish", post(publish_operation), ) .route( "/operations/{operation_id}/archive", post(archive_operation), ) .route("/operations/{operation_id}/test-runs", post(run_test)) .route( "/operations/{operation_id}/samples/input-json", post(upload_input_json), ) .route( "/operations/{operation_id}/samples/output-json", post(upload_output_json), ) .route( "/operations/{operation_id}/descriptors/proto", post(upload_proto_descriptor), ) .route( "/operations/{operation_id}/descriptors/descriptor-set", post(upload_descriptor_set), ) .route( "/operations/{operation_id}/grpc/services", get(list_grpc_services), ) .route( "/operations/{operation_id}/drafts/generate", post(generate_draft), ) .route("/operations/{operation_id}/export", get(export_operation)) .route("/agents", get(list_agents).post(create_agent)) .route( "/agents/{agent_id}", get(get_agent).patch(update_agent).delete(delete_agent), ) .route( "/agents/{agent_id}/versions/{version}", get(get_agent_version), ) .route("/agents/{agent_id}/bindings", post(save_agent_bindings)) .route("/agents/{agent_id}/publish", post(publish_agent)) .route( "/auth-profiles", get(list_auth_profiles).post(create_auth_profile), ) .route("/auth-profiles/{auth_profile_id}", get(get_auth_profile)) .route("/members", get(list_memberships)) .route( "/invitations", get(list_invitations).post(create_invitation), ) .route("/invitations/{invitation_id}", delete(delete_invitation)) .route( "/platform-api-keys", get(list_platform_api_keys).post(create_platform_api_key), ) .route( "/platform-api-keys/{key_id}/revoke", post(revoke_platform_api_key), ) .route( "/platform-api-keys/{key_id}", delete(delete_platform_api_key), ) .route("/logs", get(list_logs)) .route("/logs/{log_id}", get(get_log)) .route("/usage", get(get_usage)) .route("/usage/operations/{operation_id}", get(get_operation_usage)) .route("/usage/agents/{agent_id}", get(get_agent_usage)); let workspace_root_router = Router::new() .route("/workspaces", get(list_workspaces).post(create_workspace)) .layer(middleware::from_fn_with_state( state.clone(), require_session, )); let workspace_scoped_router = Router::new() .route( "/workspaces/{workspace_id}", get(get_workspace).patch(update_workspace), ) .nest("/workspaces/{workspace_id}", workspace_router) .layer(middleware::from_fn_with_state( state.clone(), require_workspace_session, )); let admin_router = workspace_root_router.merge(workspace_scoped_router); Router::new() .route("/health", get(crate::routes::health)) .nest( "/api/auth", Router::new() .route("/login", post(login)) .route("/logout", post(logout)) .route("/session", get(get_session)), ) .nest("/api/admin", admin_router) .with_state(state) } use crate::routes::operations::import_operation; #[cfg(test)] mod tests { use std::{ collections::BTreeMap, env, fmt, time::{SystemTime, UNIX_EPOCH}, }; use axum::{Json, Router, routing::post}; use crank_adapter_grpc::test_support as grpc_test_support; use crank_core::{ DescriptorId, ExecutionConfig, GraphqlOperationType, GraphqlTarget, GrpcTarget, HttpMethod, MembershipRole, Protocol, RestTarget, Target, ToolDescription, WorkspaceId, }; use crank_mapping::{MappingRule, MappingSet}; use crank_registry::PostgresRegistry; use crank_schema::{Schema, SchemaKind}; use serde_json::{Value, json}; use serial_test::serial; use sqlx::{Connection, Executor, PgConnection}; use tokio::net::TcpListener; use crate::{ app::build_app, auth::{AuthSettings, BootstrapAdminConfig, hash_password}, service::{AdminService, OperationPayload}, state::AppState, }; const DEFAULT_WORKSPACE_ID: &str = "ws_default"; const TEST_AUTH_EMAIL: &str = "owner@crank.local"; const TEST_AUTH_PASSWORD: &str = "test-password"; const TEST_PASSWORD_PEPPER: &str = "test-password-pepper"; const TEST_SESSION_SECRET: &str = "test-session-secret"; struct TestServer { base_url: String, shutdown: Option>, handle: Option>, } impl Drop for TestServer { fn drop(&mut self) { let shutdown = self.shutdown.take(); let handle = self.handle.take(); tokio::task::block_in_place(|| { if let Some(shutdown) = shutdown { let _ = shutdown.send(()); } if let Some(handle) = handle { let _ = tokio::runtime::Handle::current().block_on(handle); } }); } } impl fmt::Display for TestServer { fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { formatter.write_str(&self.base_url) } } impl AsRef for TestServer { fn as_ref(&self) -> &str { &self.base_url } } #[tokio::test(flavor = "multi_thread")] #[serial] async fn creates_publishes_and_tests_rest_operation() { let registry = test_registry().await; let storage_root = test_storage_root("lifecycle"); let upstream_base_url = spawn_upstream_server().await; let base_url = spawn_admin_api(build_test_app(registry, storage_root)).await; let client = authorized_client(&base_url).await; let created = client .post(format!("{base_url}/operations")) .json(&test_operation_payload( &upstream_base_url, "crm_create_lead", )) .send() .await .unwrap() .json::() .await .unwrap(); let operation_id = created["operation_id"].as_str().unwrap().to_owned(); let listed = client .get(format!("{base_url}/operations")) .send() .await .unwrap() .json::() .await .unwrap(); let published = client .post(format!("{base_url}/operations/{operation_id}/publish")) .json(&json!({ "version": 1 })) .send() .await .unwrap() .json::() .await .unwrap(); let test_run = client .post(format!("{base_url}/operations/{operation_id}/test-runs")) .json(&json!({ "version": 1, "input": { "email": "user@example.com" } })) .send() .await .unwrap() .json::() .await .unwrap(); assert_eq!(listed["items"][0]["name"], "crm_create_lead"); assert_eq!( listed["items"][0]["target_url"], format!("{upstream_base_url}/crm/leads") ); assert_eq!(listed["items"][0]["target_action"], "POST"); assert_eq!(published["published_version"], 1); assert_eq!(test_run["ok"], true); assert_eq!( test_run["request_preview"]["body"]["email"], "user@example.com" ); assert_eq!(test_run["response_preview"]["id"], "lead_123"); } #[tokio::test(flavor = "multi_thread")] #[serial] async fn updates_archives_and_deletes_operation() { let registry = test_registry().await; let storage_root = test_storage_root("operation_mutations"); let upstream_base_url = spawn_upstream_server().await; let base_url = spawn_admin_api(build_test_app(registry, storage_root)).await; let client = authorized_client(&base_url).await; let created = client .post(format!("{base_url}/operations")) .json(&test_operation_payload( &upstream_base_url, "crm_mutable_operation", )) .send() .await .unwrap() .json::() .await .unwrap(); let operation_id = created["operation_id"].as_str().unwrap().to_owned(); let listed = client .get(format!("{base_url}/operations")) .send() .await .unwrap() .json::() .await .unwrap(); assert_eq!(listed["total"], 1); assert_eq!(listed["items"][0]["category"], "sales"); assert_eq!( listed["items"][0]["target_url"], format!("{upstream_base_url}/crm/leads") ); assert_eq!(listed["items"][0]["target_action"], "POST"); let updated = client .patch(format!("{base_url}/operations/{operation_id}")) .json(&json!({ "display_name": "Create Lead Updated", "category": "marketing", "target": { "kind": "rest", "base_url": upstream_base_url, "method": "POST", "path_template": "/crm/leads", "static_headers": {} }, "input_schema": { "type": "object", "required": true, "nullable": false, "fields": { "email": { "type": "string", "required": true, "nullable": false } } }, "output_schema": { "type": "object", "required": true, "nullable": false, "fields": { "id": { "type": "string", "required": true, "nullable": false } } }, "input_mapping": { "rules": [ { "source": "$.mcp.email", "target": "$.request.body.email", "required": true } ] }, "output_mapping": { "rules": [ { "source": "$.response.body.id", "target": "$.output.id", "required": true } ] }, "execution_config": { "timeout_ms": 1000, "headers": {} }, "tool_description": { "title": "Create Lead Updated", "description": "Creates a CRM lead", "tags": ["crm"] } })) .send() .await .unwrap() .json::() .await .unwrap(); assert_eq!(updated["status"], "draft"); let detail = client .get(format!("{base_url}/operations/{operation_id}")) .send() .await .unwrap() .json::() .await .unwrap(); assert_eq!(detail["display_name"], "Create Lead Updated"); assert_eq!(detail["category"], "marketing"); let archived = client .post(format!("{base_url}/operations/{operation_id}/archive")) .send() .await .unwrap() .json::() .await .unwrap(); assert_eq!(archived["status"], "archived"); let deleted = client .delete(format!("{base_url}/operations/{operation_id}")) .send() .await .unwrap() .json::() .await .unwrap(); assert_eq!(deleted["operation_id"], operation_id); let missing = client .get(format!("{base_url}/operations/{operation_id}")) .send() .await .unwrap(); assert_eq!(missing.status(), reqwest::StatusCode::NOT_FOUND); } #[tokio::test(flavor = "multi_thread")] #[serial] async fn creates_binds_and_publishes_agent() { let registry = test_registry().await; let storage_root = test_storage_root("agent_lifecycle"); let upstream_base_url = spawn_upstream_server().await; let base_url = spawn_admin_api(build_test_app(registry, storage_root)).await; let client = authorized_client(&base_url).await; let operation = client .post(format!("{base_url}/operations")) .json(&test_operation_payload( &upstream_base_url, "crm_create_lead_agent", )) .send() .await .unwrap(); let operation = assert_success_json(operation).await; let operation_id = operation["operation_id"].as_str().unwrap().to_owned(); client .post(format!("{base_url}/operations/{operation_id}/publish")) .json(&json!({ "version": 1 })) .send() .await .unwrap(); let agent = client .post(format!("{base_url}/agents")) .json(&json!({ "slug": "sales-assistant", "display_name": "Sales Assistant", "description": "Curated sales toolset", "instructions": {}, "tool_selection_policy": {} })) .send() .await .unwrap() .json::() .await .unwrap(); let agent_id = agent["agent_id"].as_str().unwrap().to_owned(); let bindings = client .post(format!("{base_url}/agents/{agent_id}/bindings")) .json(&json!([ { "operation_id": operation_id, "operation_version": 1, "tool_name": "crm_create_lead_agent", "tool_title": "Create Lead", "enabled": true } ])) .send() .await .unwrap() .json::() .await .unwrap(); let published = client .post(format!("{base_url}/agents/{agent_id}/publish")) .json(&json!({ "version": 1 })) .send() .await .unwrap() .json::() .await .unwrap(); assert_eq!( bindings["bindings"][0]["tool_name"], "crm_create_lead_agent" ); assert_eq!(published["published_version"], 1); } #[tokio::test(flavor = "multi_thread")] #[serial] async fn updates_lists_and_deletes_agent() { let registry = test_registry().await; let storage_root = test_storage_root("agent_mutations"); let upstream_base_url = spawn_upstream_server().await; let base_url = spawn_admin_api(build_test_app(registry, storage_root)).await; let client = authorized_client(&base_url).await; let operation = client .post(format!("{base_url}/operations")) .json(&test_operation_payload( &upstream_base_url, "crm_create_lead_agents_page", )) .send() .await .unwrap(); let operation = assert_success_json(operation).await; let operation_id = operation["operation_id"].as_str().unwrap().to_owned(); client .post(format!("{base_url}/operations/{operation_id}/publish")) .json(&json!({ "version": 1 })) .send() .await .unwrap(); let created = client .post(format!("{base_url}/agents")) .json(&json!({ "slug": "support-team", "display_name": "Support Team", "description": "Support workflows", "instructions": {}, "tool_selection_policy": {} })) .send() .await .unwrap() .json::() .await .unwrap(); let agent_id = created["agent_id"].as_str().unwrap().to_owned(); client .post(format!("{base_url}/agents/{agent_id}/bindings")) .json(&json!([ { "operation_id": operation_id, "operation_version": 1, "tool_name": "crm_create_lead_agents_page", "tool_title": "Create Lead", "tool_description_override": null, "enabled": true } ])) .send() .await .unwrap(); client .post(format!("{base_url}/agents/{agent_id}/publish")) .json(&json!({ "version": 1 })) .send() .await .unwrap(); let listed = client .get(format!("{base_url}/agents")) .send() .await .unwrap() .json::() .await .unwrap(); assert_eq!(listed["items"][0]["operation_count"], 1); assert_eq!(listed["items"][0]["operation_ids"][0], operation_id); assert_eq!( listed["items"][0]["mcp_endpoint"], "/mcp/v1/default/support-team" ); assert_eq!(listed["items"][0]["status"], "published"); let updated = client .patch(format!("{base_url}/agents/{agent_id}")) .json(&json!({ "slug": "support-escalation", "display_name": "Support Escalation", "description": "Escalation workflows" })) .send() .await .unwrap() .json::() .await .unwrap(); assert_eq!(updated["agent_id"], agent_id); let detail = client .get(format!("{base_url}/agents/{agent_id}")) .send() .await .unwrap() .json::() .await .unwrap(); assert_eq!(detail["slug"], "support-escalation"); assert_eq!(detail["display_name"], "Support Escalation"); assert_eq!(detail["operation_count"], 1); assert_eq!(detail["mcp_endpoint"], "/mcp/v1/default/support-escalation"); let deleted = client .delete(format!("{base_url}/agents/{agent_id}")) .send() .await .unwrap() .json::() .await .unwrap(); assert_eq!(deleted["agent_id"], agent_id); let missing = client .get(format!("{base_url}/agents/{agent_id}")) .send() .await .unwrap(); assert_eq!(missing.status(), reqwest::StatusCode::NOT_FOUND); } #[tokio::test(flavor = "multi_thread")] #[serial] async fn manages_platform_access_resources() { let registry = test_registry().await; let storage_root = test_storage_root("platform_access"); let base_url = spawn_admin_api(build_test_app(registry, storage_root)).await; let client = authorized_client(&base_url).await; let members = client .get(format!("{base_url}/members")) .send() .await .unwrap() .json::() .await .unwrap(); let created_invitation = client .post(format!("{base_url}/invitations")) .json(&json!({ "email": "operator@example.com", "role": "operator" })) .send() .await .unwrap() .json::() .await .unwrap(); let invitation_id = created_invitation["invitation"]["invitation"]["id"] .as_str() .unwrap() .to_owned(); let invitations = client .get(format!("{base_url}/invitations")) .send() .await .unwrap() .json::() .await .unwrap(); let created_key = client .post(format!("{base_url}/platform-api-keys")) .json(&json!({ "name": "workspace-operator", "scopes": ["read", "write"] })) .send() .await .unwrap() .json::() .await .unwrap(); let key_id = created_key["api_key"]["api_key"]["id"] .as_str() .unwrap() .to_owned(); let listed_keys = client .get(format!("{base_url}/platform-api-keys")) .send() .await .unwrap() .json::() .await .unwrap(); let revoke_status = client .post(format!("{base_url}/platform-api-keys/{key_id}/revoke")) .send() .await .unwrap() .status(); let delete_invitation_status = client .delete(format!("{base_url}/invitations/{invitation_id}")) .send() .await .unwrap() .status(); let delete_key_status = client .delete(format!("{base_url}/platform-api-keys/{key_id}")) .send() .await .unwrap() .status(); assert_eq!(members["items"][0]["role"], "owner"); assert_eq!( created_invitation["invitation"]["invitation"]["status"], "pending" ); assert!( created_invitation["invite_token"] .as_str() .unwrap() .starts_with("invite_") ); assert_eq!( invitations["items"][0]["invitation"]["email"], "operator@example.com" ); assert_eq!( listed_keys["items"][0]["api_key"]["name"], "workspace-operator" ); assert!(created_key["secret"].as_str().unwrap().starts_with("crk_")); assert_eq!(revoke_status, reqwest::StatusCode::NO_CONTENT); assert_eq!(delete_invitation_status, reqwest::StatusCode::NO_CONTENT); assert_eq!(delete_key_status, reqwest::StatusCode::NO_CONTENT); } #[tokio::test(flavor = "multi_thread")] #[serial] async fn exposes_logs_and_usage_from_real_test_runs() { let registry = test_registry().await; let storage_root = test_storage_root("observability"); let upstream_base_url = spawn_upstream_server().await; let base_url = spawn_admin_api(build_test_app(registry, storage_root)).await; let client = authorized_client(&base_url).await; let created = client .post(format!("{base_url}/operations")) .json(&test_operation_payload( &upstream_base_url, "crm_observability", )) .send() .await .unwrap() .json::() .await .unwrap(); let operation_id = created["operation_id"].as_str().unwrap().to_owned(); client .post(format!("{base_url}/operations/{operation_id}/test-runs")) .json(&json!({ "version": 1, "input": { "email": "user@example.com" } })) .send() .await .unwrap(); let logs = client .get(format!("{base_url}/logs?period=7d")) .send() .await .unwrap() .json::() .await .unwrap(); let log_id = logs["items"][0]["log"]["id"].as_str().unwrap().to_owned(); let log_detail = client .get(format!("{base_url}/logs/{log_id}")) .send() .await .unwrap() .json::() .await .unwrap(); let usage = client .get(format!("{base_url}/usage?period=7d")) .send() .await .unwrap() .json::() .await .unwrap(); let operation_usage = client .get(format!( "{base_url}/usage/operations/{operation_id}?period=7d" )) .send() .await .unwrap() .json::() .await .unwrap(); assert_eq!(logs["items"][0]["log"]["source"], "admin_test_run"); assert_eq!(logs["items"][0]["operation_name"], "crm_observability"); assert_eq!(log_detail["log"]["status"], "ok"); assert_eq!(usage["summary"]["rollup"]["calls_total"], 1); assert_eq!(usage["summary"]["rollup"]["calls_ok"], 1); assert_eq!( usage["operations"][0]["operation_name"], "crm_observability" ); assert_eq!(operation_usage["rollup"]["calls_total"], 1); } #[tokio::test(flavor = "multi_thread")] #[serial] async fn creates_publishes_and_tests_graphql_operation() { let registry = test_registry().await; let storage_root = test_storage_root("graphql"); let upstream_base_url = spawn_graphql_server().await; let base_url = spawn_admin_api(build_test_app(registry, storage_root)).await; let client = authorized_client(&base_url).await; let created = client .post(format!("{base_url}/operations")) .json(&test_graphql_operation_payload( &upstream_base_url, "crm_create_lead_graphql", )) .send() .await .unwrap() .json::() .await .unwrap(); let operation_id = created["operation_id"].as_str().unwrap().to_owned(); let published = client .post(format!("{base_url}/operations/{operation_id}/publish")) .json(&json!({ "version": 1 })) .send() .await .unwrap() .json::() .await .unwrap(); let test_run = client .post(format!("{base_url}/operations/{operation_id}/test-runs")) .json(&json!({ "version": 1, "input": { "email": "user@example.com" } })) .send() .await .unwrap() .json::() .await .unwrap(); assert_eq!(published["published_version"], 1); assert_eq!(test_run["ok"], true); assert_eq!( test_run["request_preview"]["variables"]["email"], "user@example.com" ); assert_eq!(test_run["response_preview"]["id"], "lead_123"); } #[tokio::test(flavor = "multi_thread")] #[serial] async fn uploads_descriptor_set_and_lists_grpc_services() { let registry = test_registry().await; let storage_root = test_storage_root("grpc_descriptor"); let server_addr = grpc_test_support::spawn_unary_echo_server().await; let base_url = spawn_admin_api(build_test_app(registry, storage_root)).await; let client = authorized_client(&base_url).await; let created = client .post(format!("{base_url}/operations")) .json(&test_grpc_operation_payload( &server_addr, "echo_descriptor", )) .send() .await .unwrap() .json::() .await .unwrap(); let operation_id = created["operation_id"].as_str().unwrap().to_owned(); let uploaded = client .post(format!( "{base_url}/operations/{operation_id}/descriptors/descriptor-set" )) .header("x-file-name", "echo_descriptor.bin") .body(grpc_test_support::echo::FILE_DESCRIPTOR_SET.to_vec()) .send() .await .unwrap() .json::() .await .unwrap(); let services = client .get(format!( "{base_url}/operations/{operation_id}/grpc/services" )) .send() .await .unwrap() .json::() .await .unwrap(); assert!(uploaded["descriptor_id"].as_str().is_some()); assert_eq!(services["services"][0]["package"], "echo"); assert_eq!(services["services"][0]["service"], "EchoService"); assert_eq!(services["services"][0]["methods"][0]["name"], "UnaryEcho"); assert_eq!(services["services"][0]["methods"][0]["kind"], "unary"); } #[tokio::test(flavor = "multi_thread")] #[serial] async fn creates_publishes_and_tests_grpc_operation() { let registry = test_registry().await; let storage_root = test_storage_root("grpc_runtime"); let server_addr = grpc_test_support::spawn_unary_echo_server().await; let base_url = spawn_admin_api(build_test_app(registry, storage_root)).await; let client = authorized_client(&base_url).await; let created = client .post(format!("{base_url}/operations")) .json(&test_grpc_operation_payload(&server_addr, "echo_runtime")) .send() .await .unwrap() .json::() .await .unwrap(); let operation_id = created["operation_id"].as_str().unwrap().to_owned(); let published = client .post(format!("{base_url}/operations/{operation_id}/publish")) .json(&json!({ "version": 1 })) .send() .await .unwrap() .json::() .await .unwrap(); let test_run = client .post(format!("{base_url}/operations/{operation_id}/test-runs")) .json(&json!({ "version": 1, "input": { "message": "hello" } })) .send() .await .unwrap() .json::() .await .unwrap(); assert_eq!(published["published_version"], 1); assert_eq!(test_run["ok"], true); assert_eq!(test_run["request_preview"]["grpc"]["message"], "hello"); assert_eq!(test_run["response_preview"]["message"], "hello"); } #[tokio::test(flavor = "multi_thread")] #[serial] async fn manages_auth_profiles_and_yaml_upsert() { let registry = test_registry().await; let storage_root = test_storage_root("yaml"); let upstream_base_url = spawn_upstream_server().await; let base_url = spawn_admin_api(build_test_app(registry, storage_root)).await; let client = authorized_client(&base_url).await; let auth_profile = client .post(format!("{base_url}/auth-profiles")) .json(&json!({ "name": "crm-header", "kind": "api_key_header", "config": { "api_key_header": { "header_name": "X-Api-Key", "secret_ref": "secret://crm/api-key" } } })) .send() .await .unwrap() .json::() .await .unwrap(); let created = client .post(format!("{base_url}/operations")) .json(&test_operation_payload( &upstream_base_url, "crm_export_target", )) .send() .await .unwrap() .json::() .await .unwrap(); let operation_id = created["operation_id"].as_str().unwrap(); let yaml = client .get(format!("{base_url}/operations/{operation_id}/export")) .send() .await .unwrap() .text() .await .unwrap(); let imported = client .post(format!("{base_url}/operations/import?mode=upsert")) .body(yaml) .send() .await .unwrap() .json::() .await .unwrap(); assert_eq!(auth_profile["kind"], "api_key_header"); assert_eq!(imported["operation_id"], operation_id); assert_eq!(imported["version"], 2); assert_eq!(imported["import_mode"], "upsert"); } #[tokio::test(flavor = "multi_thread")] #[serial] async fn roundtrips_graphql_operation_through_yaml_upsert() { let registry = test_registry().await; let storage_root = test_storage_root("yaml_graphql"); let upstream_base_url = spawn_graphql_server().await; let base_url = spawn_admin_api(build_test_app(registry, storage_root)).await; let client = authorized_client(&base_url).await; let created = client .post(format!("{base_url}/operations")) .json(&test_graphql_operation_payload( &upstream_base_url, "crm_yaml_graphql", )) .send() .await .unwrap() .json::() .await .unwrap(); let operation_id = created["operation_id"].as_str().unwrap().to_owned(); let yaml = client .get(format!("{base_url}/operations/{operation_id}/export")) .send() .await .unwrap() .text() .await .unwrap(); let imported = client .post(format!("{base_url}/operations/import?mode=upsert")) .body(yaml) .send() .await .unwrap() .json::() .await .unwrap(); let test_run = client .post(format!("{base_url}/operations/{operation_id}/test-runs")) .json(&json!({ "version": 2, "input": { "email": "user@example.com" } })) .send() .await .unwrap() .json::() .await .unwrap(); assert_eq!(imported["operation_id"], operation_id); assert_eq!(imported["version"], 2); assert_eq!(test_run["ok"], true); assert_eq!(test_run["response_preview"]["id"], "lead_123"); } #[tokio::test(flavor = "multi_thread")] #[serial] async fn roundtrips_grpc_operation_through_yaml_upsert() { let registry = test_registry().await; let storage_root = test_storage_root("yaml_grpc"); let server_addr = grpc_test_support::spawn_unary_echo_server().await; let base_url = spawn_admin_api(build_test_app(registry, storage_root)).await; let client = authorized_client(&base_url).await; let created = client .post(format!("{base_url}/operations")) .json(&test_grpc_operation_payload(&server_addr, "echo_yaml_grpc")) .send() .await .unwrap() .json::() .await .unwrap(); let operation_id = created["operation_id"].as_str().unwrap().to_owned(); let yaml = client .get(format!("{base_url}/operations/{operation_id}/export")) .send() .await .unwrap() .text() .await .unwrap(); let imported = client .post(format!("{base_url}/operations/import?mode=upsert")) .body(yaml) .send() .await .unwrap() .json::() .await .unwrap(); let test_run = client .post(format!("{base_url}/operations/{operation_id}/test-runs")) .json(&json!({ "version": 2, "input": { "message": "hello" } })) .send() .await .unwrap() .json::() .await .unwrap(); assert_eq!(imported["operation_id"], operation_id); assert_eq!(imported["version"], 2); assert_eq!(test_run["ok"], true); assert_eq!(test_run["response_preview"]["message"], "hello"); } #[tokio::test(flavor = "multi_thread")] #[serial] async fn uploads_samples_and_generates_draft() { let registry = test_registry().await; let storage_root = test_storage_root("draft"); let upstream_base_url = spawn_upstream_server().await; let base_url = spawn_admin_api(build_test_app(registry, storage_root)).await; let client = authorized_client(&base_url).await; let created = client .post(format!("{base_url}/operations")) .json(&test_operation_payload( &upstream_base_url, "crm_draft_target", )) .send() .await .unwrap() .json::() .await .unwrap(); let operation_id = created["operation_id"].as_str().unwrap(); client .post(format!( "{base_url}/operations/{operation_id}/samples/input-json" )) .json(&json!({ "email": "user@example.com", "name": "Ada" })) .send() .await .unwrap(); client .post(format!( "{base_url}/operations/{operation_id}/samples/output-json" )) .json(&json!({ "id": "lead_123", "status": "created" })) .send() .await .unwrap(); let generated = client .post(format!( "{base_url}/operations/{operation_id}/drafts/generate" )) .json(&json!({ "sources": ["input_json_sample", "output_json_sample"] })) .send() .await .unwrap() .json::() .await .unwrap(); assert_eq!(generated["generated_draft"]["status"], "available"); assert_eq!( generated["input_schema"]["fields"]["email"]["type"], "string" ); assert_eq!( generated["output_mapping"]["rules"][0]["target"], "$.output.id" ); } fn build_test_app(registry: PostgresRegistry, storage_root: std::path::PathBuf) -> Router { build_app(AppState { service: AdminService::new(registry, storage_root, test_auth_settings()), }) } async fn spawn_upstream_server() -> String { let app = Router::new().route("/crm/leads", post(create_lead)); let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); let address = listener.local_addr().unwrap(); tokio::spawn(async move { axum::serve(listener, app).await.unwrap(); }); format!("http://{}", address) } async fn spawn_graphql_server() -> String { let app = Router::new().route("/", post(graphql_handler)); let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); let address = listener.local_addr().unwrap(); tokio::spawn(async move { axum::serve(listener, app).await.unwrap(); }); format!("http://{}", address) } async fn spawn_admin_api(app: Router) -> TestServer { let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); let address = listener.local_addr().unwrap(); let (shutdown_tx, shutdown_rx) = tokio::sync::oneshot::channel(); let handle = tokio::spawn(async move { axum::serve(listener, app) .with_graceful_shutdown(async move { let _ = shutdown_rx.await; }) .await .unwrap(); }); TestServer { base_url: format!( "http://{}/api/admin/workspaces/{}", address, DEFAULT_WORKSPACE_ID ), shutdown: Some(shutdown_tx), handle: Some(handle), } } async fn authorized_client(workspace_base_url: impl AsRef) -> reqwest::Client { let root_url = workspace_base_url .as_ref() .split("/api/admin/workspaces/") .next() .unwrap(); let client = reqwest::Client::builder() .cookie_store(true) .build() .unwrap(); client .post(format!("{root_url}/api/auth/login")) .json(&json!({ "email": TEST_AUTH_EMAIL, "password": TEST_AUTH_PASSWORD, })) .send() .await .unwrap() .error_for_status() .unwrap(); client } async fn assert_success_json(response: reqwest::Response) -> Value { let status = response.status(); let body = response.text().await.unwrap(); assert!( status.is_success(), "request failed with status {status}: {body}" ); serde_json::from_str(&body).unwrap() } async fn create_lead(Json(payload): Json) -> Json { Json(json!({ "id": "lead_123", "status": "created", "email": payload["email"] })) } async fn graphql_handler(Json(payload): Json) -> Json { let email = payload .get("variables") .and_then(|variables| variables.get("email")) .and_then(Value::as_str) .unwrap_or_default(); Json(json!({ "data": { "createLead": { "id": "lead_123", "status": "created", "email": email } } })) } async fn test_registry() -> PostgresRegistry { let database_url = env::var("TEST_DATABASE_URL") .unwrap_or_else(|_| "postgres://crank:crank@127.0.0.1:5432/crank".to_owned()); let mut admin_connection = PgConnection::connect(&database_url).await.unwrap(); let schema = format!( "test_admin_api_{}_{}", std::process::id(), SystemTime::now() .duration_since(UNIX_EPOCH) .unwrap() .as_nanos() ); admin_connection .execute(sqlx::query(&format!("create schema {schema}"))) .await .unwrap(); let registry = PostgresRegistry::connect(&format!("{database_url}?options=-csearch_path%3D{schema}")) .await .unwrap(); let password_hash = hash_password(TEST_AUTH_PASSWORD, TEST_PASSWORD_PEPPER).unwrap(); let user_id = registry .upsert_bootstrap_user(TEST_AUTH_EMAIL, "Test Owner", &password_hash) .await .unwrap(); registry .ensure_membership( &WorkspaceId::new(DEFAULT_WORKSPACE_ID), &user_id, MembershipRole::Owner, ) .await .unwrap(); registry } fn test_storage_root(name: &str) -> std::path::PathBuf { env::temp_dir().join(format!( "crank_admin_api_{name}_{}_{}", std::process::id(), SystemTime::now() .duration_since(UNIX_EPOCH) .unwrap() .as_nanos() )) } fn test_auth_settings() -> AuthSettings { AuthSettings { session_secret: TEST_SESSION_SECRET.to_owned(), password_pepper: TEST_PASSWORD_PEPPER.to_owned(), session_ttl_hours: 24, cookie_secure: false, bootstrap_admin: BootstrapAdminConfig { email: TEST_AUTH_EMAIL.to_owned(), password: TEST_AUTH_PASSWORD.to_owned(), display_name: "Test Owner".to_owned(), }, } } fn test_operation_payload(base_url: &str, name: &str) -> OperationPayload { OperationPayload { name: name.to_owned(), display_name: "Create Lead".to_owned(), category: "sales".to_owned(), protocol: Protocol::Rest, target: Target::Rest(RestTarget { base_url: base_url.to_owned(), method: HttpMethod::Post, path_template: "/crm/leads".to_owned(), static_headers: BTreeMap::new(), }), input_schema: object_schema("email"), output_schema: object_schema("id"), input_mapping: MappingSet { rules: vec![MappingRule { source: "$.mcp.email".to_owned(), target: "$.request.body.email".to_owned(), required: true, default_value: None, transform: None, condition: None, notes: None, }], }, output_mapping: MappingSet { rules: vec![MappingRule { source: "$.response.body.id".to_owned(), target: "$.output.id".to_owned(), required: true, default_value: None, transform: None, condition: None, notes: None, }], }, execution_config: ExecutionConfig { timeout_ms: 1_000, retry_policy: None, auth_profile_ref: None, headers: BTreeMap::new(), protocol_options: None, }, tool_description: ToolDescription { title: "Create Lead".to_owned(), description: "Creates a CRM lead".to_owned(), tags: vec!["crm".to_owned()], examples: Vec::new(), }, } } fn test_graphql_operation_payload(endpoint: &str, name: &str) -> OperationPayload { OperationPayload { name: name.to_owned(), display_name: "Create Lead GraphQL".to_owned(), category: "sales".to_owned(), protocol: Protocol::Graphql, target: Target::Graphql(GraphqlTarget { endpoint: endpoint.to_owned(), operation_type: GraphqlOperationType::Mutation, operation_name: "CreateLead".to_owned(), query_template: "mutation CreateLead($email: String!) { createLead(email: $email) { id status email } }" .to_owned(), response_path: "$.response.body.data.createLead".to_owned(), }), input_schema: object_schema("email"), output_schema: object_schema("id"), input_mapping: MappingSet { rules: vec![MappingRule { source: "$.mcp.email".to_owned(), target: "$.request.variables.email".to_owned(), required: true, default_value: None, transform: None, condition: None, notes: None, }], }, output_mapping: MappingSet { rules: vec![MappingRule { source: "$.response.data.id".to_owned(), target: "$.output.id".to_owned(), required: true, default_value: None, transform: None, condition: None, notes: None, }], }, execution_config: ExecutionConfig { timeout_ms: 1_000, retry_policy: None, auth_profile_ref: None, headers: BTreeMap::new(), protocol_options: None, }, tool_description: ToolDescription { title: "Create Lead GraphQL".to_owned(), description: "Creates a CRM lead through GraphQL".to_owned(), tags: vec!["crm".to_owned(), "graphql".to_owned()], examples: Vec::new(), }, } } fn test_grpc_operation_payload(server_addr: &str, name: &str) -> OperationPayload { OperationPayload { name: name.to_owned(), display_name: "Unary Echo gRPC".to_owned(), category: "support".to_owned(), protocol: Protocol::Grpc, target: Target::Grpc(GrpcTarget { server_addr: server_addr.to_owned(), package: "echo".to_owned(), service: "EchoService".to_owned(), method: "UnaryEcho".to_owned(), descriptor_ref: DescriptorId::new("desc_echo"), descriptor_set_b64: grpc_test_support::descriptor_set_b64(), }), input_schema: object_schema("message"), output_schema: object_schema("message"), input_mapping: MappingSet { rules: vec![MappingRule { source: "$.mcp.message".to_owned(), target: "$.request.grpc.message".to_owned(), required: true, default_value: None, transform: None, condition: None, notes: None, }], }, output_mapping: MappingSet { rules: vec![MappingRule { source: "$.response.data.message".to_owned(), target: "$.output.message".to_owned(), required: true, default_value: None, transform: None, condition: None, notes: None, }], }, execution_config: ExecutionConfig { timeout_ms: 1_000, retry_policy: None, auth_profile_ref: None, headers: BTreeMap::new(), protocol_options: None, }, tool_description: ToolDescription { title: "Unary Echo gRPC".to_owned(), description: "Echoes a unary gRPC payload".to_owned(), tags: vec!["grpc".to_owned()], examples: Vec::new(), }, } } fn object_schema(field_name: &str) -> Schema { Schema { kind: SchemaKind::Object, description: None, required: true, nullable: false, default_value: None, fields: BTreeMap::from([( field_name.to_owned(), Schema { kind: SchemaKind::String, description: None, required: true, nullable: false, default_value: None, fields: BTreeMap::new(), items: None, enum_values: Vec::new(), variants: Vec::new(), }, )]), items: None, enum_values: Vec::new(), variants: Vec::new(), } } }