From bf270336d99e858e3c3d8b4e3d7f9b7d65b15f6d Mon Sep 17 00:00:00 2001 From: "a.tolmachev" Date: Sun, 3 May 2026 07:40:06 +0000 Subject: [PATCH] fix: restore yaml import and async job result contracts --- .gitignore | 1 + apps/admin-api/src/app.rs | 57 +++++++++++ apps/admin-api/src/service.rs | 9 +- apps/mcp-server/src/app.rs | 92 ++++++++++-------- apps/mcp-server/src/main.rs | 148 +++++++++++++++++++++++------ crates/crank-core/src/agent.rs | 2 +- crates/crank-core/src/operation.rs | 44 ++++++++- 7 files changed, 275 insertions(+), 78 deletions(-) diff --git a/.gitignore b/.gitignore index fa8ca66..c3e88a2 100644 --- a/.gitignore +++ b/.gitignore @@ -18,3 +18,4 @@ apps/ui/vite.config.js apps/ui/vite.config.d.ts *.log __*.md +diplom.docx diff --git a/apps/admin-api/src/app.rs b/apps/admin-api/src/app.rs index f297e41..08a1f28 100644 --- a/apps/admin-api/src/app.rs +++ b/apps/admin-api/src/app.rs @@ -1940,6 +1940,63 @@ mod tests { assert!(body["error"]["context"]["poll_after_ms"].as_u64().unwrap() > 0); } + #[tokio::test(flavor = "multi_thread")] + #[serial] + async fn allows_completed_async_job_result_fetch_without_poll_delay() { + let registry = test_registry().await; + let storage_root = test_storage_root("async_job_completed_result"); + let upstream_base_url = spawn_upstream_server().await; + let base_url = spawn_admin_api(build_test_app(registry.clone(), storage_root)).await; + let client = authorized_client(&base_url).await; + + let created = client + .post(format!("{base_url}/operations")) + .json(&test_rest_async_job_operation_payload( + &upstream_base_url, + "crm_async_job_completed_result", + )) + .send() + .await + .unwrap() + .json::() + .await + .unwrap(); + let operation_id = OperationId::new(created["operation_id"].as_str().unwrap().to_owned()); + let now = OffsetDateTime::now_utc(); + let job_id = AsyncJobId::new("job_completed_result".to_owned()); + registry + .create_async_job(CreateAsyncJobRequest { + job: &AsyncJobHandle { + id: job_id.clone(), + workspace_id: WorkspaceId::new(DEFAULT_WORKSPACE_ID), + agent_id: None, + operation_id, + status: JobStatus::Completed, + progress: json!({ "pct": 100 }), + result: Some(json!({ "id": "lead_123" })), + error: None, + expires_at: Some(now + time::Duration::minutes(5)), + last_poll_at: Some(now), + created_at: now, + updated_at: now, + finished_at: Some(now), + }, + }) + .await + .unwrap(); + + let response = client + .get(format!("{base_url}/async-jobs/{}/result", job_id.as_str())) + .send() + .await + .unwrap(); + let status = response.status(); + let body = response.json::().await.unwrap(); + + assert_eq!(status, reqwest::StatusCode::OK); + assert_eq!(body["id"], "lead_123"); + } + #[tokio::test(flavor = "multi_thread")] #[serial] async fn returns_structured_context_for_missing_stream_session() { diff --git a/apps/admin-api/src/service.rs b/apps/admin-api/src/service.rs index 22c5414..fc6b47e 100644 --- a/apps/admin-api/src/service.rs +++ b/apps/admin-api/src/service.rs @@ -1719,9 +1719,12 @@ impl AdminService { ) })?; ensure_async_job_workspace(&job, workspace_id)?; - let job = self - .touch_async_job_poll_if_ready(workspace_id, job) - .await?; + let job = if matches!(job.status, JobStatus::Completed) { + job + } else { + self.touch_async_job_poll_if_ready(workspace_id, job) + .await? + }; match job.status { JobStatus::Completed => Ok(job.result.unwrap_or(Value::Null)), diff --git a/apps/mcp-server/src/app.rs b/apps/mcp-server/src/app.rs index 367332e..f6c63b3 100644 --- a/apps/mcp-server/src/app.rs +++ b/apps/mcp-server/src/app.rs @@ -1416,28 +1416,32 @@ async fn handle_async_job_status_call( ); }; - let now = OffsetDateTime::now_utc(); - let poll_after_ms = streaming.poll_interval_ms.unwrap_or(1_000); - let remaining_delay_ms = job.remaining_poll_delay_ms(now, poll_after_ms); - if remaining_delay_ms > 0 { - return tool_error_response( - message, - response_mode, - &session.protocol_version, - "async_job_poll_rate_limited", - format!( - "async job {} must wait before the next poll", - control.job_id - ), - Some(json!({ - "poll_after_ms": remaining_delay_ms, - })), - ); - } + let job = if matches!(job.status, JobStatus::Completed) { + job + } else { + let now = OffsetDateTime::now_utc(); + let poll_after_ms = streaming.poll_interval_ms.unwrap_or(1_000); + let remaining_delay_ms = job.remaining_poll_delay_ms(now, poll_after_ms); + if remaining_delay_ms > 0 { + return tool_error_response( + message, + response_mode, + &session.protocol_version, + "async_job_poll_rate_limited", + format!( + "async job {} must wait before the next poll", + control.job_id + ), + Some(json!({ + "poll_after_ms": remaining_delay_ms, + })), + ); + } - let job = match state.registry.touch_async_job_poll(&job.id, &now).await { - Ok(job) => job, - Err(error) => return internal_jsonrpc_error(message, error), + match state.registry.touch_async_job_poll(&job.id, &now).await { + Ok(job) => job, + Err(error) => return internal_jsonrpc_error(message, error), + } }; success_tool_response( @@ -1516,28 +1520,32 @@ async fn handle_async_job_result_call( ); }; - let now = OffsetDateTime::now_utc(); - let poll_after_ms = streaming.poll_interval_ms.unwrap_or(1_000); - let remaining_delay_ms = job.remaining_poll_delay_ms(now, poll_after_ms); - if remaining_delay_ms > 0 { - return tool_error_response( - message, - response_mode, - &session.protocol_version, - "async_job_poll_rate_limited", - format!( - "async job {} must wait before the next poll", - control.job_id - ), - Some(json!({ - "poll_after_ms": remaining_delay_ms, - })), - ); - } + let job = if matches!(job.status, JobStatus::Completed) { + job + } else { + let now = OffsetDateTime::now_utc(); + let poll_after_ms = streaming.poll_interval_ms.unwrap_or(1_000); + let remaining_delay_ms = job.remaining_poll_delay_ms(now, poll_after_ms); + if remaining_delay_ms > 0 { + return tool_error_response( + message, + response_mode, + &session.protocol_version, + "async_job_poll_rate_limited", + format!( + "async job {} must wait before the next poll", + control.job_id + ), + Some(json!({ + "poll_after_ms": remaining_delay_ms, + })), + ); + } - let job = match state.registry.touch_async_job_poll(&job.id, &now).await { - Ok(job) => job, - Err(error) => return internal_jsonrpc_error(message, error), + match state.registry.touch_async_job_poll(&job.id, &now).await { + Ok(job) => job, + Err(error) => return internal_jsonrpc_error(message, error), + } }; match job.status { diff --git a/apps/mcp-server/src/main.rs b/apps/mcp-server/src/main.rs index 641e924..ad4fc94 100644 --- a/apps/mcp-server/src/main.rs +++ b/apps/mcp-server/src/main.rs @@ -135,15 +135,16 @@ mod tests { use crank_adapter_grpc::test_support as grpc_test_support; use crank_core::{ Agent, AgentId, AgentOperationBinding, AgentStatus, AgentVersion, AggregationMode, - DescriptorId, ExecutionConfig, ExecutionMode, GraphqlOperationType, GraphqlTarget, - GrpcTarget, HttpMethod, Operation, OperationId, OperationStatus, PlatformApiKey, - PlatformApiKeyId, PlatformApiKeyScope, PlatformApiKeyStatus, Protocol, RestTarget, - StreamingConfig, Target, ToolDescription, ToolFamilyConfig, TransportBehavior, WorkspaceId, + AsyncJobHandle, AsyncJobId, DescriptorId, ExecutionConfig, ExecutionMode, + GraphqlOperationType, GraphqlTarget, GrpcTarget, HttpMethod, JobStatus, Operation, + OperationId, OperationStatus, PlatformApiKey, PlatformApiKeyId, PlatformApiKeyScope, + PlatformApiKeyStatus, Protocol, RestTarget, StreamingConfig, Target, ToolDescription, + ToolFamilyConfig, TransportBehavior, WorkspaceId, }; use crank_mapping::{MappingRule, MappingSet}; use crank_registry::{ - CreateAgentRequest, CreatePlatformApiKeyRequest, ListInvocationLogsQuery, PostgresRegistry, - PublishAgentRequest, PublishRequest, + CreateAgentRequest, CreateAsyncJobRequest, CreatePlatformApiKeyRequest, + ListInvocationLogsQuery, PostgresRegistry, PublishAgentRequest, PublishRequest, }; use crank_runtime::{ RequestRateLimitConfig, RequestRateLimiter, RuntimeExecutor, SecretCrypto, @@ -1759,6 +1760,29 @@ mod tests { ) .await; + let now = OffsetDateTime::now_utc(); + let job_id = AsyncJobId::new("job_async_rate".to_owned()); + registry + .create_async_job(CreateAsyncJobRequest { + job: &AsyncJobHandle { + id: job_id.clone(), + workspace_id: test_workspace_id(), + agent_id: Some(AgentId::new("agent_sales-async-rate")), + operation_id: operation.id.clone(), + status: JobStatus::Running, + progress: json!({ "pct": 25 }), + result: None, + error: None, + expires_at: Some(now + time::Duration::minutes(5)), + last_poll_at: None, + created_at: now, + updated_at: now, + finished_at: None, + }, + }) + .await + .unwrap(); + let base_url = spawn_mcp_server(build_test_app( registry, Duration::from_millis(0), @@ -1769,29 +1793,6 @@ mod tests { let mcp_url = agent_mcp_url(&base_url, "sales-async-rate"); let initialized_session = initialize_session(&client, &mcp_url, &api_key).await; - let start_response = post_jsonrpc( - &client, - &mcp_url, - &api_key, - Some(&initialized_session), - json!({ - "jsonrpc": "2.0", - "id": 1, - "method": "tools/call", - "params": { - "name": "crm_async_rate_start", - "arguments": { - "email": "user@example.com" - } - } - }), - ) - .await; - let job_id = start_response["result"]["structuredContent"]["job_id"] - .as_str() - .unwrap() - .to_owned(); - let first_status = post_jsonrpc( &client, &mcp_url, @@ -1804,7 +1805,7 @@ mod tests { "params": { "name": "crm_async_rate_status", "arguments": { - "job_id": job_id + "job_id": job_id.as_str() } } }), @@ -1824,7 +1825,7 @@ mod tests { "params": { "name": "crm_async_rate_status", "arguments": { - "job_id": job_id + "job_id": job_id.as_str() } } }), @@ -1842,6 +1843,91 @@ mod tests { assert!((1..=250).contains(&poll_after_ms)); } + #[tokio::test(flavor = "multi_thread")] + async fn allows_completed_async_job_result_without_poll_delay() { + let registry = test_registry().await; + let upstream_base_url = spawn_upstream_server().await; + let operation = test_rest_async_job_operation(&upstream_base_url, "crm_async_completed"); + + registry + .create_operation(&test_workspace_id(), &operation, Some("alice")) + .await + .unwrap(); + registry + .publish_operation(PublishRequest { + workspace_id: &test_workspace_id(), + operation_id: &operation.id, + version: 1, + published_at: &OffsetDateTime::parse("2026-03-26T10:00:00Z", &Rfc3339).unwrap(), + published_by: Some("alice"), + }) + .await + .unwrap(); + publish_agent_for_operation(®istry, &operation, "sales-async-completed").await; + let api_key = create_platform_api_key( + ®istry, + "mcp-async-completed", + &[PlatformApiKeyScope::Read, PlatformApiKeyScope::Write], + ) + .await; + + let now = OffsetDateTime::now_utc(); + let job_id = AsyncJobId::new("job_async_completed".to_owned()); + registry + .create_async_job(CreateAsyncJobRequest { + job: &AsyncJobHandle { + id: job_id.clone(), + workspace_id: test_workspace_id(), + agent_id: Some(AgentId::new("agent_sales-async-completed")), + operation_id: operation.id.clone(), + status: JobStatus::Completed, + progress: json!({ "pct": 100 }), + result: Some(json!({ "id": "lead_123" })), + error: None, + expires_at: Some(now + time::Duration::minutes(5)), + last_poll_at: Some(now), + created_at: now, + updated_at: now, + finished_at: Some(now), + }, + }) + .await + .unwrap(); + + let base_url = spawn_mcp_server(build_test_app( + registry, + Duration::from_millis(0), + Some("https://crank.example.com".to_owned()), + )) + .await; + let client = reqwest::Client::new(); + let mcp_url = agent_mcp_url(&base_url, "sales-async-completed"); + let initialized_session = initialize_session(&client, &mcp_url, &api_key).await; + + let result_response = post_jsonrpc( + &client, + &mcp_url, + &api_key, + Some(&initialized_session), + json!({ + "jsonrpc": "2.0", + "id": 1, + "method": "tools/call", + "params": { + "name": "crm_async_completed_result", + "arguments": { + "job_id": job_id.as_str() + } + } + }), + ) + .await; + assert_eq!( + result_response["result"]["structuredContent"]["id"], + "lead_123" + ); + } + #[tokio::test] async fn rejects_rapid_initialize_requests_with_429() { let registry = test_registry().await; diff --git a/crates/crank-core/src/agent.rs b/crates/crank-core/src/agent.rs index 7b4857e..35b772a 100644 --- a/crates/crank-core/src/agent.rs +++ b/crates/crank-core/src/agent.rs @@ -26,7 +26,7 @@ pub struct Agent { pub created_at: OffsetDateTime, #[serde(with = "time::serde::rfc3339")] pub updated_at: OffsetDateTime, - #[serde(with = "time::serde::rfc3339::option")] + #[serde(default, with = "time::serde::rfc3339::option")] pub published_at: Option, } diff --git a/crates/crank-core/src/operation.rs b/crates/crank-core/src/operation.rs index 091cde5..7d7cad7 100644 --- a/crates/crank-core/src/operation.rs +++ b/crates/crank-core/src/operation.rs @@ -236,7 +236,7 @@ pub struct Operation { pub created_at: OffsetDateTime, #[serde(with = "time::serde::rfc3339")] pub updated_at: OffsetDateTime, - #[serde(skip_serializing_if = "Option::is_none")] + #[serde(default, skip_serializing_if = "Option::is_none")] #[serde(with = "time::serde::rfc3339::option")] pub published_at: Option, } @@ -517,4 +517,46 @@ mod tests { assert!(yaml.contains("export_mode: portable")); assert_eq!(restored, operation); } + + #[test] + fn operation_yaml_deserializes_without_optional_published_at() { + let yaml = r#" +id: op_01 +name: crm_create_lead +display_name: Create Lead +category: sales +protocol: rest +status: draft +version: 1 +target: + kind: rest + base_url: https://api.example.com + method: POST + path_template: /v1/leads + static_headers: {} +input_schema: + type: object +output_schema: + type: object +input_mapping: + rules: [] +output_mapping: + rules: [] +execution_config: + timeout_ms: 10000 + headers: {} +tool_description: + title: Create CRM lead + description: Creates a new lead. + tags: [] + examples: [] +created_at: 2026-03-25T08:00:00Z +updated_at: 2026-03-25T08:10:00Z +"#; + + let restored: Operation = + serde_yaml::from_str(yaml).unwrap(); + + assert_eq!(restored.published_at, None); + } }