registry: add sqlx checks for agent and operation reads

This commit is contained in:
a.tolmachev
2026-04-12 21:07:30 +00:00
parent 6a909feb64
commit 934c872753
10 changed files with 1176 additions and 165 deletions
+273 -123
View File
@@ -594,56 +594,6 @@ fn map_invocation_log_record(row: &PgRow) -> Result<InvocationLogRecord, Registr
})
}
fn map_agent_summary(row: &PgRow) -> Result<AgentSummary, RegistryError> {
Ok(AgentSummary {
id: AgentId::new(row.try_get::<String, _>("id")?),
workspace_id: WorkspaceId::new(row.try_get::<String, _>("workspace_id")?),
slug: row.try_get("slug")?,
display_name: row.try_get("display_name")?,
description: row.try_get("description")?,
status: deserialize_enum_text(&row.try_get::<String, _>("status")?, "status")?,
current_draft_version: from_db_version(
row.try_get("current_draft_version")?,
"current_draft_version",
)?,
latest_published_version: row
.try_get::<Option<i32>, _>("latest_published_version")?
.map(|value| from_db_version(value, "latest_published_version"))
.transpose()?,
created_at: row.try_get("created_at")?,
updated_at: row.try_get("updated_at")?,
published_at: row.try_get("published_at")?,
})
}
fn map_operation_summary(row: &PgRow) -> Result<OperationSummary, RegistryError> {
let target: Target = deserialize_json_value(row.try_get::<Json<Value>, _>("target_json")?.0)?;
let (target_url, target_action) = target_summary(&target);
Ok(OperationSummary {
id: OperationId::new(row.try_get::<String, _>("id")?),
workspace_id: WorkspaceId::new(row.try_get::<String, _>("workspace_id")?),
name: row.try_get("name")?,
display_name: row.try_get("display_name")?,
category: row.try_get("category")?,
protocol: deserialize_enum_text(&row.try_get::<String, _>("protocol")?, "protocol")?,
target_url,
target_action,
status: deserialize_enum_text(&row.try_get::<String, _>("status")?, "status")?,
current_draft_version: from_db_version(
row.try_get("current_draft_version")?,
"current_draft_version",
)?,
latest_published_version: row
.try_get::<Option<i32>, _>("latest_published_version")?
.map(|value| from_db_version(value, "latest_published_version"))
.transpose()?,
created_at: row.try_get("created_at")?,
updated_at: row.try_get("updated_at")?,
published_at: row.try_get("published_at")?,
})
}
fn map_operation_usage_summary(row: &PgRow) -> Result<OperationUsageSummary, RegistryError> {
Ok(OperationUsageSummary {
operation_id: OperationId::new(row.try_get::<String, _>("operation_id")?),
@@ -716,62 +666,35 @@ fn map_operation_agent_ref(row: &PgRow) -> Result<OperationAgentRef, RegistryErr
}
fn map_operation_version_record(row: &PgRow) -> Result<OperationVersionRecord, RegistryError> {
let operation_id = OperationId::new(row.try_get::<String, _>("id")?);
let version = from_db_version(row.try_get("version")?, "version")?;
let status = deserialize_enum_text(&row.try_get::<String, _>("status")?, "status")?;
Ok(OperationVersionRecord {
operation_id: operation_id.clone(),
workspace_id: WorkspaceId::new(row.try_get::<String, _>("workspace_id")?),
version,
status,
change_note: row.try_get("change_note")?,
created_at: row.try_get("created_at")?,
created_by: row.try_get("created_by")?,
snapshot: RegistryOperation {
id: operation_id,
name: row.try_get("name")?,
display_name: row.try_get("display_name")?,
category: row.try_get("category")?,
protocol: deserialize_enum_text(&row.try_get::<String, _>("protocol")?, "protocol")?,
status,
version,
target: deserialize_json_value(row.try_get::<Json<Value>, _>("target_json")?.0)?,
input_schema: deserialize_json_value(
row.try_get::<Json<Value>, _>("input_schema_json")?.0,
)?,
output_schema: deserialize_json_value(
row.try_get::<Json<Value>, _>("output_schema_json")?.0,
)?,
input_mapping: deserialize_json_value(
row.try_get::<Json<Value>, _>("input_mapping_json")?.0,
)?,
output_mapping: deserialize_json_value(
row.try_get::<Json<Value>, _>("output_mapping_json")?.0,
)?,
execution_config: deserialize_json_value(
row.try_get::<Json<Value>, _>("execution_config_json")?.0,
)?,
tool_description: deserialize_json_value(
row.try_get::<Json<Value>, _>("tool_description_json")?.0,
)?,
samples: row
.try_get::<Option<Json<Value>>, _>("samples_json")?
.map(|value| deserialize_json_value(value.0))
.transpose()?,
generated_draft: row
.try_get::<Option<Json<Value>>, _>("generated_draft_json")?
.map(|value| deserialize_json_value(value.0))
.transpose()?,
config_export: row
.try_get::<Option<Json<Value>>, _>("config_export_json")?
.map(|value| deserialize_json_value(value.0))
.transpose()?,
created_at: row.try_get("operation_created_at")?,
updated_at: row.try_get("operation_updated_at")?,
published_at: row.try_get("operation_published_at")?,
},
})
build_operation_version_record(
row.try_get("id")?,
row.try_get("workspace_id")?,
row.try_get("name")?,
row.try_get("display_name")?,
row.try_get("category")?,
row.try_get("protocol")?,
row.try_get("operation_created_at")?,
row.try_get("operation_updated_at")?,
row.try_get("operation_published_at")?,
row.try_get("version")?,
row.try_get("status")?,
row.try_get::<Json<Value>, _>("target_json")?.0,
row.try_get::<Json<Value>, _>("input_schema_json")?.0,
row.try_get::<Json<Value>, _>("output_schema_json")?.0,
row.try_get::<Json<Value>, _>("input_mapping_json")?.0,
row.try_get::<Json<Value>, _>("output_mapping_json")?.0,
row.try_get::<Json<Value>, _>("execution_config_json")?.0,
row.try_get::<Json<Value>, _>("tool_description_json")?.0,
row.try_get::<Option<Json<Value>>, _>("samples_json")?
.map(|value| value.0),
row.try_get::<Option<Json<Value>>, _>("generated_draft_json")?
.map(|value| value.0),
row.try_get::<Option<Json<Value>>, _>("config_export_json")?
.map(|value| value.0),
row.try_get("change_note")?,
row.try_get("created_at")?,
row.try_get("created_by")?,
)
}
fn map_auth_profile(row: &PgRow) -> Result<AuthProfile, RegistryError> {
@@ -797,30 +720,168 @@ fn map_agent_binding(row: &PgRow) -> Result<AgentOperationBinding, RegistryError
})
}
fn map_agent_version_record(
#[allow(clippy::too_many_arguments)]
fn build_agent_summary(
id: String,
workspace_id: String,
slug: String,
display_name: String,
description: String,
status: String,
current_draft_version: i32,
latest_published_version: Option<i32>,
created_at: String,
updated_at: String,
published_at: Option<String>,
) -> Result<AgentSummary, RegistryError> {
Ok(AgentSummary {
id: AgentId::new(id),
workspace_id: WorkspaceId::new(workspace_id),
slug,
display_name,
description,
status: deserialize_enum_text(&status, "status")?,
current_draft_version: from_db_version(current_draft_version, "current_draft_version")?,
latest_published_version: latest_published_version
.map(|value| from_db_version(value, "latest_published_version"))
.transpose()?,
created_at,
updated_at,
published_at,
})
}
#[allow(clippy::too_many_arguments)]
fn build_operation_summary(
id: String,
workspace_id: String,
name: String,
display_name: String,
category: String,
protocol: String,
target_json: Value,
status: String,
current_draft_version: i32,
latest_published_version: Option<i32>,
created_at: String,
updated_at: String,
published_at: Option<String>,
) -> Result<OperationSummary, RegistryError> {
let target: Target = deserialize_json_value(target_json)?;
let (target_url, target_action) = target_summary(&target);
Ok(OperationSummary {
id: OperationId::new(id),
workspace_id: WorkspaceId::new(workspace_id),
name,
display_name,
category,
protocol: deserialize_enum_text(&protocol, "protocol")?,
target_url,
target_action,
status: deserialize_enum_text(&status, "status")?,
current_draft_version: from_db_version(current_draft_version, "current_draft_version")?,
latest_published_version: latest_published_version
.map(|value| from_db_version(value, "latest_published_version"))
.transpose()?,
created_at,
updated_at,
published_at,
})
}
fn build_agent_version_record(
summary: &AgentSummary,
row: &PgRow,
version: i32,
status: String,
instructions_json: Value,
tool_selection_policy_json: Value,
created_at: String,
bindings: Vec<AgentOperationBinding>,
) -> Result<AgentVersionRecord, RegistryError> {
let version = from_db_version(row.try_get("version")?, "version")?;
let status = deserialize_enum_text(&row.try_get::<String, _>("status")?, "status")?;
let version = from_db_version(version, "version")?;
let status = deserialize_enum_text(&status, "status")?;
Ok(AgentVersionRecord {
agent_id: summary.id.clone(),
workspace_id: summary.workspace_id.clone(),
version,
status,
created_at: row.try_get("created_at")?,
created_at: created_at.clone(),
bindings,
snapshot: AgentVersion {
agent_id: summary.id.clone(),
version,
status,
instructions: row.try_get::<Json<Value>, _>("instructions_json")?.0,
tool_selection_policy: row
.try_get::<Json<Value>, _>("tool_selection_policy_json")?
.0,
created_at: row.try_get("created_at")?,
instructions: instructions_json,
tool_selection_policy: tool_selection_policy_json,
created_at,
},
})
}
#[allow(clippy::too_many_arguments)]
fn build_operation_version_record(
id: String,
workspace_id: String,
name: String,
display_name: String,
category: String,
protocol: String,
operation_created_at: String,
operation_updated_at: String,
operation_published_at: Option<String>,
version: i32,
status: String,
target_json: Value,
input_schema_json: Value,
output_schema_json: Value,
input_mapping_json: Value,
output_mapping_json: Value,
execution_config_json: Value,
tool_description_json: Value,
samples_json: Option<Value>,
generated_draft_json: Option<Value>,
config_export_json: Option<Value>,
change_note: Option<String>,
created_at: String,
created_by: Option<String>,
) -> Result<OperationVersionRecord, RegistryError> {
let operation_id = OperationId::new(id);
let version = from_db_version(version, "version")?;
let status = deserialize_enum_text(&status, "status")?;
Ok(OperationVersionRecord {
operation_id: operation_id.clone(),
workspace_id: WorkspaceId::new(workspace_id),
version,
status,
change_note,
created_at: created_at.clone(),
created_by,
snapshot: RegistryOperation {
id: operation_id,
name,
display_name,
category,
protocol: deserialize_enum_text(&protocol, "protocol")?,
status,
version,
target: deserialize_json_value(target_json)?,
input_schema: deserialize_json_value(input_schema_json)?,
output_schema: deserialize_json_value(output_schema_json)?,
input_mapping: deserialize_json_value(input_mapping_json)?,
output_mapping: deserialize_json_value(output_mapping_json)?,
execution_config: deserialize_json_value(execution_config_json)?,
tool_description: deserialize_json_value(tool_description_json)?,
samples: samples_json.map(deserialize_json_value).transpose()?,
generated_draft: generated_draft_json
.map(deserialize_json_value)
.transpose()?,
config_export: config_export_json.map(deserialize_json_value).transpose()?,
created_at: operation_created_at,
updated_at: operation_updated_at,
published_at: operation_published_at,
},
})
}
@@ -1046,12 +1107,13 @@ mod tests {
};
use crank_core::{
ApiKeyHeaderAuthConfig, AsyncJobHandle, AuthConfig, AuthKind, AuthProfile, ConfigExport,
ExecutionConfig, ExportMode, GeneratedDraft, GeneratedDraftStatus, HttpMethod, JobStatus,
MembershipRole, OperationId, OperationStatus, PlatformApiKey, PlatformApiKeyId,
PlatformApiKeyScope, PlatformApiKeyStatus, Protocol, RestTarget, RetryPolicy, Samples,
SecretId, StreamSession, StreamSessionId, StreamStatus, Target, ToolDescription,
ToolExample, User, UserId, Workspace, WorkspaceId,
AgentId, AgentOperationBinding, AgentStatus, AgentVersion, ApiKeyHeaderAuthConfig,
AsyncJobHandle, AuthConfig, AuthKind, AuthProfile, ConfigExport, ExecutionConfig,
ExportMode, GeneratedDraft, GeneratedDraftStatus, HttpMethod, JobStatus, MembershipRole,
OperationId, OperationStatus, PlatformApiKey, PlatformApiKeyId, PlatformApiKeyScope,
PlatformApiKeyStatus, Protocol, RestTarget, RetryPolicy, Samples, SecretId, StreamSession,
StreamSessionId, StreamStatus, Target, ToolDescription, ToolExample, User, UserId,
Workspace, WorkspaceId,
};
use crank_mapping::{MappingRule, MappingSet};
use crank_schema::{Schema, SchemaKind};
@@ -1061,7 +1123,7 @@ mod tests {
use crate::{
PostgresRegistry, RegistryError,
model::{
AsyncJobFilter, CreateAsyncJobRequest, CreatePlatformApiKeyRequest,
AsyncJobFilter, CreateAgentRequest, CreateAsyncJobRequest, CreatePlatformApiKeyRequest,
CreateStreamSessionRequest, CreateVersionRequest, CreateWorkspaceRequest,
CreateYamlImportJobRequest, DescriptorKind, DescriptorMetadata,
OperationSampleMetadata, PlatformApiKeyRecord, PublishRequest, RegistryOperation,
@@ -1524,6 +1586,61 @@ mod tests {
database.cleanup().await;
}
#[tokio::test]
async fn manages_agent_read_paths() {
let database = TestDatabase::new().await;
let registry = database.registry().await;
let operation = test_operation("op_agent_01", 1, OperationStatus::Draft);
let agent = test_agent("agent_01", AgentStatus::Draft);
let version = test_agent_version(&agent.id, 1, AgentStatus::Draft);
let bindings = vec![AgentOperationBinding {
operation_id: operation.id.clone(),
operation_version: operation.version,
tool_name: "create_lead".to_owned(),
tool_title: "Create lead".to_owned(),
tool_description_override: Some("Creates CRM lead".to_owned()),
enabled: true,
}];
registry
.create_operation(&test_workspace_id(), &operation, None)
.await
.unwrap();
registry
.create_agent(CreateAgentRequest {
agent: &agent,
version: &version,
bindings: &bindings,
})
.await
.unwrap();
let listed = registry.list_agents(&test_workspace_id()).await.unwrap();
let summary = registry
.get_agent_summary(&test_workspace_id(), &agent.id)
.await
.unwrap()
.unwrap();
let loaded_version = registry
.get_agent_version(&test_workspace_id(), &agent.id, version.version)
.await
.unwrap()
.unwrap();
assert_eq!(listed, vec![summary.clone()]);
assert_eq!(summary.id, agent.id);
assert_eq!(summary.slug, agent.slug);
assert_eq!(summary.display_name, agent.display_name);
assert_eq!(summary.status, agent.status);
assert_eq!(loaded_version.agent_id, agent.id);
assert_eq!(loaded_version.version, version.version);
assert_eq!(loaded_version.status, version.status);
assert_eq!(loaded_version.snapshot, version);
assert_eq!(loaded_version.bindings, bindings);
database.cleanup().await;
}
#[tokio::test]
async fn manages_stream_sessions_with_transitions_and_cleanup() {
let database = TestDatabase::new().await;
@@ -1981,6 +2098,39 @@ mod tests {
}
}
fn test_agent(id: &str, status: AgentStatus) -> crank_core::Agent {
crank_core::Agent {
id: AgentId::new(id),
workspace_id: test_workspace_id(),
slug: format!("{id}_slug"),
display_name: format!("Display {id}"),
description: format!("Description {id}"),
status,
current_draft_version: 1,
latest_published_version: None,
created_at: "2026-03-25T11:58:00Z".to_owned(),
updated_at: "2026-03-25T12:00:00Z".to_owned(),
published_at: None,
}
}
fn test_agent_version(agent_id: &AgentId, version: u32, status: AgentStatus) -> AgentVersion {
AgentVersion {
agent_id: agent_id.clone(),
version,
status,
instructions: json!({
"system": "triage tickets",
"guardrails": ["don't mutate state"]
}),
tool_selection_policy: json!({
"mode": "allow_list",
"max_tools": 8
}),
created_at: "2026-03-25T12:00:00Z".to_owned(),
}
}
fn test_async_job(id: &str, operation_id: &str, status: JobStatus) -> AsyncJobHandle {
AsyncJobHandle {
id: crank_core::AsyncJobId::new(id),