From b2c5e28cbabb4a139e90722fc762cf9d1ef4f865 Mon Sep 17 00:00:00 2001 From: "a.tolmachev" Date: Mon, 4 May 2026 09:04:27 +0000 Subject: [PATCH] cache: add graphql query response caching --- TASKS.md | 11 +- apps/admin-api/src/service.rs | 81 ++++++++++++- crates/crank-core/src/cache.rs | 1 + crates/crank-runtime/src/cache.rs | 3 + crates/crank-runtime/src/executor.rs | 164 ++++++++++++++++++++++++++- docs/admin-api.md | 3 +- docs/implementation-plan.md | 6 +- docs/mcp-interface.md | 6 +- docs/runtime-config.md | 8 ++ 9 files changed, 271 insertions(+), 12 deletions(-) diff --git a/TASKS.md b/TASKS.md index 4e8f5dc..f80c6bb 100644 --- a/TASKS.md +++ b/TASKS.md @@ -38,6 +38,10 @@ Implementation slices: - ключи response cache изолируются минимум по `workspace + agent + operation + request fingerprint`; - ключи coordination state изолируются минимум по `workspace + agent + cache scope`; - при отсутствии внешнего cache все эти контуры продолжают работать через in-memory fallback. +7. Довести response cache не только до Community baseline, но и до первых коммерчески ценных read-only hot paths: + - `REST GET` + - `GraphQL query` + при сохранении той же изоляции и opt-in policy. DoD: - без `Valkey/Redis` Community и self-hosted deployment остаются рабочими; @@ -46,6 +50,10 @@ DoD: - `Enterprise` может подключать cache layer на своей инфраструктуре, включая кластерный вариант; - runtime/docs/manifests не делают cache обязательным для базового запуска. - cache boundaries и key namespaces документированы так, чтобы разные агенты и рабочие области не пересекались друг с другом. +- response cache покрывает первый коммерчески значимый protocol set без перехода к небезопасному "cache everything": + - `REST GET` + - `GraphQL query` +- дальнейшее расширение на `gRPC unary` возможно только для явно read-only / cacheable операций. Verification: - targeted runtime/admin/mcp tests for cache contracts and fallback mode; @@ -64,8 +72,9 @@ Progress: - runtime executor now performs actual response cache reads/writes for eligible `REST GET` operations with keys isolated by `workspace + agent + operation + request fingerprint` - response cache keys now also include `operation version`, so a new published version does not reuse stale cached payloads from the previous one - mcp-server published tool catalog now uses shared coordination snapshots across instances, while preserving the same `workspace + agent` isolation model + - response cache runtime/admin contracts now also cover `GraphQL query` as the first commercial protocol extension beyond the `REST GET` baseline - pending: - - preserve the same isolation model while extending response cache beyond the first `REST GET` hot path + - preserve the same isolation model while extending response cache to the next monetizable read-only path: explicit `gRPC unary` read operations - wire replay guard into the first real token or nonce flow once commercial machine-token issuance stops being a `Community` stub - wire coordination state into additional runtime or transport paths beyond the shared published catalog snapshot diff --git a/apps/admin-api/src/service.rs b/apps/admin-api/src/service.rs index b26b65b..470d680 100644 --- a/apps/admin-api/src/service.rs +++ b/apps/admin-api/src/service.rs @@ -5063,9 +5063,12 @@ fn validate_response_cache_policy( match target { Target::Rest(rest_target) if rest_target.method == crank_core::HttpMethod::Get => {} + Target::Graphql(graphql_target) + if graphql_target.operation_type == crank_core::GraphqlOperationType::Query => {} _ => { return Err(ApiError::validation_with_context( - "response cache is supported only for REST GET operations".to_owned(), + "response cache is supported only for REST GET and GraphQL query operations" + .to_owned(), json!({ "field": "execution_config.response_cache", }), @@ -5085,6 +5088,82 @@ fn validate_response_cache_policy( Ok(()) } +#[cfg(test)] +mod tests { + use std::collections::BTreeMap; + + use crank_core::{ + ExecutionConfig, GraphqlOperationType, GraphqlTarget, HttpMethod, ResponseCachePolicy, + RestTarget, Target, + }; + + use super::validate_response_cache_policy; + + fn cacheable_execution_config() -> ExecutionConfig { + ExecutionConfig { + timeout_ms: 1_000, + retry_policy: None, + response_cache: Some(ResponseCachePolicy { ttl_ms: 5_000 }), + auth_profile_ref: None, + headers: BTreeMap::new(), + protocol_options: None, + streaming: None, + } + } + + #[test] + fn accepts_response_cache_for_graphql_query() { + let target = Target::Graphql(GraphqlTarget { + endpoint: "http://example.invalid/graphql".to_owned(), + operation_type: GraphqlOperationType::Query, + operation_name: "LookupLead".to_owned(), + query_template: + "query LookupLead($email: String!) { lookupLead(email: $email) { id } }".to_owned(), + response_path: "$.response.body.data.lookupLead".to_owned(), + }); + + let result = validate_response_cache_policy(&target, &cacheable_execution_config()); + + assert!(result.is_ok()); + } + + #[test] + fn rejects_response_cache_for_graphql_mutation() { + let target = Target::Graphql(GraphqlTarget { + endpoint: "http://example.invalid/graphql".to_owned(), + operation_type: GraphqlOperationType::Mutation, + operation_name: "CreateLead".to_owned(), + query_template: + "mutation CreateLead($email: String!) { createLead(email: $email) { id } }" + .to_owned(), + response_path: "$.response.body.data.createLead".to_owned(), + }); + + let error = + validate_response_cache_policy(&target, &cacheable_execution_config()).unwrap_err(); + + assert!(matches!(error, crate::error::ApiError::Validation { .. })); + assert_eq!( + error.to_string(), + "response cache is supported only for REST GET and GraphQL query operations" + ); + } + + #[test] + fn accepts_response_cache_for_rest_get() { + let target = Target::Rest(RestTarget { + base_url: "http://example.invalid".to_owned(), + method: HttpMethod::Get, + path_template: "/catalog".to_owned(), + static_headers: BTreeMap::new(), + }); + + let result = validate_response_cache_policy(&target, &cacheable_execution_config()); + + assert!(result.is_ok()); + } +} + fn enrich_operation_summary( summary: OperationSummary, usage_summary: OperationUsageSummaryView, diff --git a/crates/crank-core/src/cache.rs b/crates/crank-core/src/cache.rs index d3dd98b..c3badd8 100644 --- a/crates/crank-core/src/cache.rs +++ b/crates/crank-core/src/cache.rs @@ -65,6 +65,7 @@ pub struct CachedResponse { pub status: u16, pub headers: Vec, pub body: Vec, + pub data: Vec, } #[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize, Deserialize)] diff --git a/crates/crank-runtime/src/cache.rs b/crates/crank-runtime/src/cache.rs index ebc7888..a7dfcfe 100644 --- a/crates/crank-runtime/src/cache.rs +++ b/crates/crank-runtime/src/cache.rs @@ -679,6 +679,7 @@ mod tests { value: "application/json".to_owned(), }], body: br#"{"ok":true}"#.to_vec(), + data: br#"{"ok":true}"#.to_vec(), }; store @@ -825,6 +826,7 @@ mod tests { status: 200, headers: vec![], body: b"ok".to_vec(), + data: br#"null"#.to_vec(), }, Duration::from_secs(5), ) @@ -868,6 +870,7 @@ mod tests { value: "hit".to_owned(), }], body: br#"{"queued":true}"#.to_vec(), + data: br#"null"#.to_vec(), }; let encoded = serde_json::to_vec(&response).unwrap(); diff --git a/crates/crank-runtime/src/executor.rs b/crates/crank-runtime/src/executor.rs index 1910717..86792f7 100644 --- a/crates/crank-runtime/src/executor.rs +++ b/crates/crank-runtime/src/executor.rs @@ -757,10 +757,12 @@ fn log_runtime_event( } fn response_cache_ttl(operation: &RuntimeOperation) -> Option { - if !matches!( - &operation.target, - Target::Rest(target) if target.method == HttpMethod::Get - ) { + let is_cacheable_protocol = match &operation.target { + Target::Rest(target) => target.method == HttpMethod::Get, + Target::Graphql(target) => target.operation_type == crank_core::GraphqlOperationType::Query, + _ => false, + }; + if !is_cacheable_protocol { return None; } @@ -784,6 +786,8 @@ fn response_cache_key( "query_params": prepared_request.query_params, "headers": prepared_request.headers, "body": prepared_request.body, + "variables": prepared_request.variables, + "grpc": prepared_request.grpc, }); let fingerprint_bytes = serde_json::to_vec(&fingerprint).ok()?; let fingerprint_hash = URL_SAFE_NO_PAD.encode(Sha256::digest(fingerprint_bytes)); @@ -811,6 +815,7 @@ fn cached_response_from_adapter(adapter_response: &AdapterResponse) -> Option Result { let body: Value = serde_json::from_slice(&cached_response.body).map_err(|error| error.to_string())?; + let data: Value = + serde_json::from_slice(&cached_response.data).map_err(|error| error.to_string())?; Ok(AdapterResponse { status_code: cached_response.status, headers: cached_response @@ -827,7 +834,7 @@ fn adapter_response_from_cached( .map(|header| (header.name, header.value)) .collect(), body, - data: Value::Null, + data, }) } @@ -1089,6 +1096,40 @@ mod tests { assert_eq!(request_count.load(Ordering::SeqCst), 2); } + #[tokio::test] + async fn caches_graphql_query_responses_within_agent_scope() { + let request_count = Arc::new(AtomicUsize::new(0)); + let endpoint = spawn_cached_graphql_server(Arc::clone(&request_count)).await; + let executor = RuntimeExecutor::new() + .with_response_cache_store(Arc::new(InMemoryResponseCacheStore::default())); + let operation = test_cached_graphql_query_operation(&endpoint); + let first_context = RuntimeRequestContext::from_request_id("req_cache_graphql_1") + .with_response_cache_scope("ws_cache", "agent_sales"); + let second_context = RuntimeRequestContext::from_request_id("req_cache_graphql_2") + .with_response_cache_scope("ws_cache", "agent_sales"); + + let first_output = executor + .execute_with_context( + &operation, + &json!({ "email": "user@example.com" }), + Some(&first_context), + ) + .await + .unwrap(); + let second_output = executor + .execute_with_context( + &operation, + &json!({ "email": "user@example.com" }), + Some(&second_context), + ) + .await + .unwrap(); + + assert_eq!(first_output, json!({ "id": "lead_123" })); + assert_eq!(second_output, json!({ "id": "lead_123" })); + assert_eq!(request_count.load(Ordering::SeqCst), 1); + } + #[tokio::test] async fn emits_runtime_tracing_with_request_context() { let base_url = spawn_runtime_server().await; @@ -1607,6 +1648,41 @@ mod tests { format!("http://{}", address) } + async fn spawn_cached_graphql_server(request_count: Arc) -> String { + let app = Router::new().route( + "/", + post(move |Json(payload): Json| { + let request_count = Arc::clone(&request_count); + async move { + let count = request_count.fetch_add(1, Ordering::SeqCst) + 1; + let email = payload + .get("variables") + .and_then(|variables| variables.get("email")) + .and_then(Value::as_str) + .unwrap_or_default(); + + Json(json!({ + "data": { + "lookupLead": { + "id": "lead_123", + "email": email, + "request_count": count, + } + } + })) + } + }), + ); + 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_context_server() -> String { let app = Router::new().route("/", post(graphql_context_handler)); let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); @@ -2090,6 +2166,84 @@ mod tests { }) } + fn test_cached_graphql_query_operation(endpoint: &str) -> RuntimeOperation { + RuntimeOperation::from(Operation { + id: OperationId::new("op_graphql_cached_runtime"), + name: "crm_lookup_lead_graphql".to_owned(), + display_name: "Lookup Lead GraphQL".to_owned(), + category: "sales".to_owned(), + protocol: Protocol::Graphql, + security_level: OperationSecurityLevel::Standard, + status: OperationStatus::Published, + version: 1, + target: Target::Graphql(GraphqlTarget { + endpoint: endpoint.to_owned(), + operation_type: GraphqlOperationType::Query, + operation_name: "LookupLead".to_owned(), + query_template: + "query LookupLead($email: String!) { lookupLead(email: $email) { id email request_count } }" + .to_owned(), + response_path: "$.response.body.data.lookupLead".to_owned(), + }), + input_schema: object_schema("email", SchemaKind::String), + output_schema: object_schema("id", SchemaKind::String), + 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, + response_cache: Some(crank_core::ResponseCachePolicy { ttl_ms: 5_000 }), + auth_profile_ref: None, + headers: BTreeMap::new(), + protocol_options: None, + streaming: None, + }, + tool_description: ToolDescription { + title: "Lookup Lead GraphQL".to_owned(), + description: "Reads a CRM lead through GraphQL".to_owned(), + tags: vec!["crm".to_owned(), "graphql".to_owned()], + examples: vec![ToolExample { + input: json!({ "email": "user@example.com" }), + }], + }, + samples: Some(Samples::default()), + generated_draft: Some(GeneratedDraft { + status: GeneratedDraftStatus::Available, + source_types: vec!["input_json".to_owned()], + generated_at: Some("2026-03-25T20:00:00Z".to_owned()), + input_schema_generated: true, + output_schema_generated: true, + input_mapping_generated: true, + output_mapping_generated: true, + warnings: Vec::new(), + }), + config_export: None, + created_at: timestamp("2026-03-25T20:00:00Z"), + updated_at: timestamp("2026-03-25T20:00:00Z"), + published_at: Some(timestamp("2026-03-25T20:00:00Z")), + }) + } + fn test_grpc_operation(server_addr: &str) -> RuntimeOperation { RuntimeOperation::from(Operation { id: OperationId::new("op_grpc_runtime"), diff --git a/docs/admin-api.md b/docs/admin-api.md index c351f01..cc45441 100644 --- a/docs/admin-api.md +++ b/docs/admin-api.md @@ -141,9 +141,10 @@ - `execution_config.response_cache` допускается только для: - `REST` - `GET` + - или `GraphQL query` - операций без `auth_profile_ref` - `execution_config.response_cache` не означает глобальный shared cache на все вызовы системы: - - response cache должен быть изолирован минимум по `workspace + agent + operation + request fingerprint` + - response cache должен быть изолирован минимум по `workspace + agent + operation + operation version + request fingerprint` - попытка создать, обновить или импортировать операцию с неподдерживаемым `protocol` или `security_level` должна завершаться `validation_error` еще на стороне `admin-api`, а не только скрываться в UI. ### 5.4. Samples and descriptors diff --git a/docs/implementation-plan.md b/docs/implementation-plan.md index b81a567..90efcb5 100644 --- a/docs/implementation-plan.md +++ b/docs/implementation-plan.md @@ -184,7 +184,7 @@ ### Цель Подготовить optional cache/coordination layer так, чтобы платформа могла использовать `Valkey/Redis`, -но не зависела от него для базового запуска. +но не зависела от него для базового запуска, и при этом закрыть первые коммерчески ценные cache hot paths. ### DoD @@ -198,6 +198,10 @@ - response cache по умолчанию изолировался по `workspace + agent + operation + operation version + request fingerprint`; - внешний cache backend можно было безопасно шарить между несколькими рабочими областями и агентами. - shared coordination cache уже применяется не только для rate limiting, но и для multi-instance snapshots published MCP catalogs. +- response cache уже покрывает: + - `REST GET` как Community baseline; + - `GraphQL query` как первый коммерчески ценный read-only protocol path. +- дальнейшее расширение response cache идет только по явно обоснованным read-only сценариям, начиная с `gRPC unary`, а не как общий cache для всех протоколов подряд. ## 14. Этап 11. Live staging and demo readiness diff --git a/docs/mcp-interface.md b/docs/mcp-interface.md index 4eabdd7..2dce7e1 100644 --- a/docs/mcp-interface.md +++ b/docs/mcp-interface.md @@ -237,9 +237,9 @@ Crank должен корректно работать, если MCP client де - `POST` и `GET` transport semantics соответствуют MCP spec `2025-06-18`; - endpoint определяется парой `workspace + agent`; - одна published operation = один MCP tool внутри agent; -- если операция явно включает `execution_config.response_cache`, первый поддержанный runtime path ограничен: - - `REST` - - `GET` +- если операция явно включает `execution_config.response_cache`, текущий поддержанный runtime path ограничен: + - `REST GET` + - `GraphQL query` - без `auth_profile_ref` - cache keys изолируются минимум по `workspace + agent + operation + operation version + request fingerprint` - streaming operations публикуются как bounded tools или tool families; diff --git a/docs/runtime-config.md b/docs/runtime-config.md index afb2ba6..c37559b 100644 --- a/docs/runtime-config.md +++ b/docs/runtime-config.md @@ -183,6 +183,14 @@ Demo/deployment: Второй контур хранит только кэшируемые ответы операций. +Текущий безопасный runtime scope для response cache: + +- `REST GET`; +- `GraphQL query`. + +Это не означает автоматическое кэширование всех protocol families. Следующим кандидатом на расширение +может быть только явно read-only `gRPC unary`. + Эти контуры не должны смешивать ключи друг с другом. ### Cache key isolation