fix: restore yaml import and async job result contracts

This commit is contained in:
a.tolmachev
2026-05-03 07:40:06 +00:00
parent 420e9416e5
commit bf270336d9
7 changed files with 275 additions and 78 deletions
+1
View File
@@ -18,3 +18,4 @@ apps/ui/vite.config.js
apps/ui/vite.config.d.ts
*.log
__*.md
diplom.docx
+57
View File
@@ -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::<Value>()
.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::<Value>().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() {
+6 -3
View File
@@ -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)),
+50 -42
View File
@@ -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 {
+117 -31
View File
@@ -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(&registry, &operation, "sales-async-completed").await;
let api_key = create_platform_api_key(
&registry,
"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;
+1 -1
View File
@@ -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<OffsetDateTime>,
}
+43 -1
View File
@@ -236,7 +236,7 @@ pub struct Operation<TSchema, TMapping> {
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<OffsetDateTime>,
}
@@ -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_json::Value, serde_json::Value> =
serde_yaml::from_str(yaml).unwrap();
assert_eq!(restored.published_at, None);
}
}