feat: add observability api foundation

This commit is contained in:
a.tolmachev
2026-03-30 00:00:13 +03:00
parent 0a1680f24e
commit be9ee95cbe
21 changed files with 1535 additions and 94 deletions
+1
View File
@@ -45,3 +45,4 @@ define_id!(UserId);
define_id!(AgentId);
define_id!(InvitationId);
define_id!(PlatformApiKeyId);
define_id!(InvocationLogId);
+6 -2
View File
@@ -2,6 +2,7 @@ pub mod access;
pub mod agent;
pub mod auth;
pub mod ids;
pub mod observability;
pub mod operation;
pub mod protocol;
pub mod workspace;
@@ -16,8 +17,11 @@ pub use auth::{
BearerAuthConfig, SecretRef,
};
pub use ids::{
AgentId, AuthProfileId, DescriptorId, InvitationId, OperationId, PlatformApiKeyId, SampleId,
ToolId, UserId, WorkspaceId,
AgentId, AuthProfileId, DescriptorId, InvitationId, InvocationLogId, OperationId,
PlatformApiKeyId, SampleId, ToolId, UserId, WorkspaceId,
};
pub use observability::{
InvocationLevel, InvocationLog, InvocationSource, InvocationStatus, UsagePeriod, UsageRollup,
};
pub use operation::{
ConfigExport, ExecutionConfig, GeneratedDraft, GeneratedDraftStatus, GraphqlTarget,
+81
View File
@@ -0,0 +1,81 @@
use serde::{Deserialize, Serialize};
use serde_json::Value;
use crate::{AgentId, OperationId, WorkspaceId};
#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum InvocationSource {
AdminTestRun,
AgentToolCall,
}
#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum InvocationLevel {
Debug,
Info,
Warn,
Error,
}
#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum InvocationStatus {
Ok,
Error,
}
#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize, Deserialize)]
pub enum UsagePeriod {
#[serde(rename = "30m")]
Last30Minutes,
#[serde(rename = "1h")]
LastHour,
#[serde(rename = "6h")]
Last6Hours,
#[serde(rename = "24h")]
Last24Hours,
#[serde(rename = "7d")]
Last7Days,
#[serde(rename = "30d")]
Last30Days,
#[serde(rename = "90d")]
Last90Days,
#[serde(rename = "this_month")]
ThisMonth,
}
#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
pub struct InvocationLog {
pub id: crate::ids::InvocationLogId,
pub workspace_id: WorkspaceId,
pub agent_id: Option<AgentId>,
pub operation_id: OperationId,
pub source: InvocationSource,
pub level: InvocationLevel,
pub status: InvocationStatus,
pub tool_name: String,
pub message: String,
pub request_id: Option<String>,
pub status_code: Option<u16>,
pub duration_ms: u64,
pub error_kind: Option<String>,
pub request_preview: Value,
pub response_preview: Value,
pub created_at: String,
}
#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
pub struct UsageRollup {
pub workspace_id: WorkspaceId,
pub agent_id: Option<AgentId>,
pub operation_id: Option<OperationId>,
pub period: UsagePeriod,
pub calls_total: u64,
pub calls_ok: u64,
pub calls_error: u64,
pub p50_ms: u64,
pub p95_ms: u64,
pub p99_ms: u64,
}
+2
View File
@@ -14,6 +14,8 @@ pub enum RegistryError {
InvitationNotFound { invitation_id: String },
#[error("platform api key {key_id} was not found")]
PlatformApiKeyNotFound { key_id: String },
#[error("invocation log {log_id} was not found")]
InvocationLogNotFound { log_id: String },
#[error("agent {agent_id} was not found")]
AgentNotFound { agent_id: String },
#[error("published agent {workspace_slug}/{agent_slug} was not found")]
+8 -6
View File
@@ -6,12 +6,14 @@ mod postgres;
pub use error::RegistryError;
pub use model::{
AgentSummary, AgentVersionRecord, CreateAgentRequest, CreateInvitationRequest,
CreatePlatformApiKeyRequest, CreateVersionRequest, CreateWorkspaceRequest,
CreateYamlImportJobRequest, DescriptorKind, DescriptorMetadata, InvitationRecord,
MembershipRecord, OperationSampleMetadata, OperationSummary, OperationVersionRecord,
PlatformApiKeyRecord, PublishAgentRequest, PublishRequest, PublishedAgentTool,
RegistryOperation, SampleKind, SaveAgentBindingsRequest, SaveAuthProfileRequest,
SaveDescriptorMetadataRequest, SaveSampleMetadataRequest, UpdateWorkspaceRequest,
CreateInvocationLogRequest, CreatePlatformApiKeyRequest, CreateVersionRequest,
CreateWorkspaceRequest, CreateYamlImportJobRequest, DescriptorKind, DescriptorMetadata,
InvitationRecord, InvocationLogRecord, ListInvocationLogsQuery, MembershipRecord,
OperationSampleMetadata, OperationSummary, OperationVersionRecord, PlatformApiKeyRecord,
PublishAgentRequest, PublishRequest, PublishedAgentTool, RegistryOperation, SampleKind,
SaveAgentBindingsRequest, SaveAuthProfileRequest, SaveDescriptorMetadataRequest,
SaveSampleMetadataRequest, UpdateWorkspaceRequest, UsageAgentBreakdown, UsageBucket,
UsageOperationBreakdown, UsageQuery, UsageRollupRecord, UsageSummary, UsageTimelinePoint,
WorkspaceRecord, YamlImportJob, YamlImportJobCompletion, YamlImportJobId, YamlImportJobStatus,
};
pub use postgres::PostgresRegistry;
+59
View File
@@ -385,5 +385,64 @@ pub async fn apply_postgres(pool: &PgPool) -> Result<(), sqlx::Error> {
.execute(pool)
.await?;
query(
"create table if not exists invocation_logs (
id text primary key,
workspace_id text not null references workspaces(id) on delete cascade,
agent_id text null references agents(id) on delete set null,
operation_id text not null references operations(id) on delete cascade,
source text not null,
level text not null,
status text not null,
tool_name text not null,
message text not null,
request_id text null,
status_code integer null,
duration_ms bigint not null,
error_kind text null,
request_preview_json jsonb not null,
response_preview_json jsonb not null,
created_at timestamptz not null
)",
)
.execute(pool)
.await?;
query(
"create index if not exists invocation_logs_workspace_created_idx on invocation_logs(workspace_id, created_at desc)",
)
.execute(pool)
.await?;
query(
"create index if not exists invocation_logs_workspace_operation_created_idx on invocation_logs(workspace_id, operation_id, created_at desc)",
)
.execute(pool)
.await?;
query(
"create index if not exists invocation_logs_workspace_agent_created_idx on invocation_logs(workspace_id, agent_id, created_at desc)",
)
.execute(pool)
.await?;
query(
"create table if not exists usage_rollups (
workspace_id text not null references workspaces(id) on delete cascade,
agent_id text null references agents(id) on delete cascade,
operation_id text null references operations(id) on delete cascade,
period text not null,
calls_total bigint not null,
calls_ok bigint not null,
calls_error bigint not null,
p50_ms bigint not null,
p95_ms bigint not null,
p99_ms bigint not null,
updated_at timestamptz not null
)",
)
.execute(pool)
.await?;
Ok(())
}
+89 -2
View File
@@ -1,7 +1,8 @@
use crank_core::{
Agent, AgentId, AgentOperationBinding, AgentStatus, AgentVersion, AuthProfile, DescriptorId,
ExportMode, InvitationToken, MembershipRole, Operation, OperationId, OperationStatus,
PlatformApiKey, Protocol, SampleId, User, Workspace, WorkspaceId,
ExportMode, InvitationToken, InvocationLevel, InvocationLog, InvocationSource, MembershipRole,
Operation, OperationId, OperationStatus, PlatformApiKey, Protocol, SampleId, UsagePeriod,
UsageRollup, User, Workspace, WorkspaceId,
};
use crank_mapping::MappingSet;
use crank_schema::Schema;
@@ -64,6 +65,52 @@ pub struct PlatformApiKeyRecord {
pub api_key: PlatformApiKey,
}
#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
pub struct InvocationLogRecord {
pub log: InvocationLog,
pub operation_name: String,
pub operation_display_name: String,
pub agent_slug: Option<String>,
pub agent_display_name: Option<String>,
}
#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
pub struct UsageRollupRecord {
pub rollup: UsageRollup,
}
#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
pub struct UsageTimelinePoint {
pub bucket_start: String,
pub calls_ok: u64,
pub calls_error: u64,
}
#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
pub struct UsageOperationBreakdown {
pub operation_id: OperationId,
pub operation_name: String,
pub operation_display_name: String,
pub protocol: Protocol,
pub calls_total: u64,
pub calls_error: u64,
pub p50_ms: u64,
pub p95_ms: u64,
pub p99_ms: u64,
}
#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
pub struct UsageAgentBreakdown {
pub agent_id: AgentId,
pub agent_slug: String,
pub agent_display_name: String,
pub calls_total: u64,
pub calls_error: u64,
pub p50_ms: u64,
pub p95_ms: u64,
pub p99_ms: u64,
}
#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
pub struct AgentSummary {
pub id: AgentId,
@@ -200,6 +247,46 @@ pub struct YamlImportJobCompletion {
pub finished_at: String,
}
#[derive(Clone, Debug, PartialEq)]
pub struct CreateInvocationLogRequest<'a> {
pub log: &'a InvocationLog,
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub struct ListInvocationLogsQuery<'a> {
pub workspace_id: &'a WorkspaceId,
pub level: Option<InvocationLevel>,
pub search_text: Option<&'a str>,
pub source: Option<InvocationSource>,
pub operation_id: Option<&'a OperationId>,
pub agent_id: Option<&'a AgentId>,
pub created_after: Option<&'a str>,
pub limit: u32,
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub enum UsageBucket {
Hour,
Day,
Week,
Month,
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub struct UsageQuery<'a> {
pub workspace_id: &'a WorkspaceId,
pub period: UsagePeriod,
pub source: Option<InvocationSource>,
pub created_after: &'a str,
pub bucket: UsageBucket,
}
#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
pub struct UsageSummary {
pub rollup: UsageRollup,
pub success_rate: f64,
}
#[derive(Clone, Debug, PartialEq)]
pub struct CreateVersionRequest<'a> {
pub workspace_id: &'a WorkspaceId,
+532 -8
View File
@@ -1,7 +1,7 @@
use crank_core::{
AgentId, AgentOperationBinding, AgentStatus, AgentVersion, AuthProfile, InvitationId,
InvitationToken, OperationId, OperationStatus, PlatformApiKey, PlatformApiKeyId,
PlatformApiKeyStatus, User, UserId, Workspace, WorkspaceId,
InvitationToken, InvocationLog, InvocationLogId, OperationId, OperationStatus, PlatformApiKey,
PlatformApiKeyId, PlatformApiKeyStatus, UsageRollup, User, UserId, Workspace, WorkspaceId,
};
use serde::{Serialize, de::DeserializeOwned};
use serde_json::Value;
@@ -16,12 +16,14 @@ use crate::{
migrations,
model::{
AgentSummary, AgentVersionRecord, CreateAgentRequest, CreateInvitationRequest,
CreatePlatformApiKeyRequest, CreateVersionRequest, CreateWorkspaceRequest,
CreateYamlImportJobRequest, DescriptorMetadata, InvitationRecord, MembershipRecord,
OperationSampleMetadata, OperationSummary, OperationVersionRecord, PlatformApiKeyRecord,
PublishAgentRequest, PublishRequest, PublishedAgentTool, RegistryOperation,
SaveAgentBindingsRequest, SaveAuthProfileRequest, SaveDescriptorMetadataRequest,
SaveSampleMetadataRequest, UpdateWorkspaceRequest, WorkspaceRecord, YamlImportJob,
CreateInvocationLogRequest, CreatePlatformApiKeyRequest, CreateVersionRequest,
CreateWorkspaceRequest, CreateYamlImportJobRequest, DescriptorMetadata, InvitationRecord,
InvocationLogRecord, ListInvocationLogsQuery, MembershipRecord, OperationSampleMetadata,
OperationSummary, OperationVersionRecord, PlatformApiKeyRecord, PublishAgentRequest,
PublishRequest, PublishedAgentTool, RegistryOperation, SaveAgentBindingsRequest,
SaveAuthProfileRequest, SaveDescriptorMetadataRequest, SaveSampleMetadataRequest,
UpdateWorkspaceRequest, UsageAgentBreakdown, UsageOperationBreakdown, UsageQuery,
UsageRollupRecord, UsageSummary, UsageTimelinePoint, WorkspaceRecord, YamlImportJob,
YamlImportJobCompletion, YamlImportJobId, YamlImportJobStatus,
},
};
@@ -289,6 +291,449 @@ impl PostgresRegistry {
Ok(())
}
pub async fn create_invocation_log(
&self,
request: CreateInvocationLogRequest<'_>,
) -> Result<(), RegistryError> {
sqlx::query(
"insert into invocation_logs (
id,
workspace_id,
agent_id,
operation_id,
source,
level,
status,
tool_name,
message,
request_id,
status_code,
duration_ms,
error_kind,
request_preview_json,
response_preview_json,
created_at
) values (
$1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, $13, $14, $15, $16::timestamptz
)",
)
.bind(request.log.id.as_str())
.bind(request.log.workspace_id.as_str())
.bind(request.log.agent_id.as_ref().map(|value| value.as_str()))
.bind(request.log.operation_id.as_str())
.bind(serialize_enum_text(&request.log.source, "source")?)
.bind(serialize_enum_text(&request.log.level, "level")?)
.bind(serialize_enum_text(&request.log.status, "status")?)
.bind(&request.log.tool_name)
.bind(&request.log.message)
.bind(&request.log.request_id)
.bind(request.log.status_code.map(i32::from))
.bind(i64::try_from(request.log.duration_ms).map_err(|_| {
RegistryError::InvalidNumericValue {
field: "duration_ms",
value: i64::MAX,
}
})?)
.bind(&request.log.error_kind)
.bind(Json(request.log.request_preview.clone()))
.bind(Json(request.log.response_preview.clone()))
.bind(&request.log.created_at)
.execute(&self.pool)
.await?;
Ok(())
}
pub async fn list_invocation_logs(
&self,
query: ListInvocationLogsQuery<'_>,
) -> Result<Vec<InvocationLogRecord>, RegistryError> {
let rows = sqlx::query(
"select
l.id,
l.workspace_id,
l.agent_id,
l.operation_id,
l.source,
l.level,
l.status,
l.tool_name,
l.message,
l.request_id,
l.status_code,
l.duration_ms,
l.error_kind,
l.request_preview_json,
l.response_preview_json,
to_char(l.created_at at time zone 'UTC', 'YYYY-MM-DD\"T\"HH24:MI:SS.MS\"Z\"') as created_at,
o.name as operation_name,
o.display_name as operation_display_name,
a.slug as agent_slug,
a.display_name as agent_display_name
from invocation_logs l
join operations o on o.id = l.operation_id
left join agents a on a.id = l.agent_id
where l.workspace_id = $1
and ($2::text is null or l.level = $2)
and ($3::text is null or l.source = $3)
and ($4::text is null or l.operation_id = $4)
and ($5::text is null or l.agent_id = $5)
and ($6::timestamptz is null or l.created_at >= $6::timestamptz)
and (
$7::text is null
or l.tool_name ilike '%' || $7 || '%'
or l.message ilike '%' || $7 || '%'
or o.name ilike '%' || $7 || '%'
or o.display_name ilike '%' || $7 || '%'
)
order by l.created_at desc
limit $8",
)
.bind(query.workspace_id.as_str())
.bind(
query
.level
.as_ref()
.map(|value| serialize_enum_text(value, "level"))
.transpose()?,
)
.bind(
query
.source
.as_ref()
.map(|value| serialize_enum_text(value, "source"))
.transpose()?,
)
.bind(query.operation_id.map(|value| value.as_str()))
.bind(query.agent_id.map(|value| value.as_str()))
.bind(query.created_after)
.bind(query.search_text)
.bind(i64::from(query.limit))
.fetch_all(&self.pool)
.await?;
rows.iter().map(map_invocation_log_record).collect()
}
pub async fn get_invocation_log(
&self,
workspace_id: &WorkspaceId,
log_id: &InvocationLogId,
) -> Result<Option<InvocationLogRecord>, RegistryError> {
let row = sqlx::query(
"select
l.id,
l.workspace_id,
l.agent_id,
l.operation_id,
l.source,
l.level,
l.status,
l.tool_name,
l.message,
l.request_id,
l.status_code,
l.duration_ms,
l.error_kind,
l.request_preview_json,
l.response_preview_json,
to_char(l.created_at at time zone 'UTC', 'YYYY-MM-DD\"T\"HH24:MI:SS.MS\"Z\"') as created_at,
o.name as operation_name,
o.display_name as operation_display_name,
a.slug as agent_slug,
a.display_name as agent_display_name
from invocation_logs l
join operations o on o.id = l.operation_id
left join agents a on a.id = l.agent_id
where l.workspace_id = $1 and l.id = $2",
)
.bind(workspace_id.as_str())
.bind(log_id.as_str())
.fetch_optional(&self.pool)
.await?;
row.as_ref().map(map_invocation_log_record).transpose()
}
pub async fn summarize_usage(
&self,
query: UsageQuery<'_>,
) -> Result<UsageSummary, RegistryError> {
let row = sqlx::query(
"select
count(*)::bigint as calls_total,
count(*) filter (where status = 'ok')::bigint as calls_ok,
count(*) filter (where status = 'error')::bigint as calls_error,
coalesce(percentile_cont(0.5) within group (order by duration_ms), 0)::bigint as p50_ms,
coalesce(percentile_cont(0.95) within group (order by duration_ms), 0)::bigint as p95_ms,
coalesce(percentile_cont(0.99) within group (order by duration_ms), 0)::bigint as p99_ms
from invocation_logs
where workspace_id = $1
and created_at >= $2::timestamptz
and ($3::text is null or source = $3)",
)
.bind(query.workspace_id.as_str())
.bind(query.created_after)
.bind(
query
.source
.as_ref()
.map(|value| serialize_enum_text(value, "source"))
.transpose()?,
)
.fetch_one(&self.pool)
.await?;
let calls_total = row.try_get::<i64, _>("calls_total")?;
let calls_ok = row.try_get::<i64, _>("calls_ok")?;
let calls_error = row.try_get::<i64, _>("calls_error")?;
let rollup = UsageRollup {
workspace_id: query.workspace_id.clone(),
agent_id: None,
operation_id: None,
period: query.period,
calls_total: to_u64(calls_total, "calls_total")?,
calls_ok: to_u64(calls_ok, "calls_ok")?,
calls_error: to_u64(calls_error, "calls_error")?,
p50_ms: to_u64(row.try_get::<i64, _>("p50_ms")?, "p50_ms")?,
p95_ms: to_u64(row.try_get::<i64, _>("p95_ms")?, "p95_ms")?,
p99_ms: to_u64(row.try_get::<i64, _>("p99_ms")?, "p99_ms")?,
};
let success_rate = if rollup.calls_total == 0 {
0.0
} else {
(rollup.calls_ok as f64 / rollup.calls_total as f64) * 100.0
};
Ok(UsageSummary {
rollup,
success_rate,
})
}
pub async fn list_usage_timeline(
&self,
query: UsageQuery<'_>,
) -> Result<Vec<UsageTimelinePoint>, RegistryError> {
let bucket = usage_bucket_sql(query.bucket);
let sql = format!(
"select
to_char(date_trunc('{bucket}', created_at at time zone 'UTC'), 'YYYY-MM-DD\"T\"HH24:MI:SS\"Z\"') as bucket_start,
count(*) filter (where status = 'ok')::bigint as calls_ok,
count(*) filter (where status = 'error')::bigint as calls_error
from invocation_logs
where workspace_id = $1
and created_at >= $2::timestamptz
and ($3::text is null or source = $3)
group by 1
order by 1 asc"
);
let rows = sqlx::query(&sql)
.bind(query.workspace_id.as_str())
.bind(query.created_after)
.bind(
query
.source
.as_ref()
.map(|value| serialize_enum_text(value, "source"))
.transpose()?,
)
.fetch_all(&self.pool)
.await?;
rows.iter()
.map(|row| {
Ok(UsageTimelinePoint {
bucket_start: row.try_get("bucket_start")?,
calls_ok: to_u64(row.try_get::<i64, _>("calls_ok")?, "calls_ok")?,
calls_error: to_u64(row.try_get::<i64, _>("calls_error")?, "calls_error")?,
})
})
.collect()
}
pub async fn list_usage_by_operation(
&self,
query: UsageQuery<'_>,
) -> Result<Vec<UsageOperationBreakdown>, RegistryError> {
let rows = sqlx::query(
"select
o.id as operation_id,
o.name as operation_name,
o.display_name as operation_display_name,
o.protocol,
count(*)::bigint as calls_total,
count(*) filter (where l.status = 'error')::bigint as calls_error,
coalesce(percentile_cont(0.5) within group (order by l.duration_ms), 0)::bigint as p50_ms,
coalesce(percentile_cont(0.95) within group (order by l.duration_ms), 0)::bigint as p95_ms,
coalesce(percentile_cont(0.99) within group (order by l.duration_ms), 0)::bigint as p99_ms
from invocation_logs l
join operations o on o.id = l.operation_id
where l.workspace_id = $1
and l.created_at >= $2::timestamptz
and ($3::text is null or l.source = $3)
group by o.id, o.name, o.display_name, o.protocol
order by calls_total desc, o.name asc",
)
.bind(query.workspace_id.as_str())
.bind(query.created_after)
.bind(
query
.source
.as_ref()
.map(|value| serialize_enum_text(value, "source"))
.transpose()?,
)
.fetch_all(&self.pool)
.await?;
rows.iter().map(map_usage_operation_breakdown).collect()
}
pub async fn get_usage_for_operation(
&self,
query: UsageQuery<'_>,
operation_id: &OperationId,
) -> Result<Option<UsageRollupRecord>, RegistryError> {
let row = sqlx::query(
"select
count(*)::bigint as calls_total,
count(*) filter (where status = 'ok')::bigint as calls_ok,
count(*) filter (where status = 'error')::bigint as calls_error,
coalesce(percentile_cont(0.5) within group (order by duration_ms), 0)::bigint as p50_ms,
coalesce(percentile_cont(0.95) within group (order by duration_ms), 0)::bigint as p95_ms,
coalesce(percentile_cont(0.99) within group (order by duration_ms), 0)::bigint as p99_ms
from invocation_logs
where workspace_id = $1
and operation_id = $2
and created_at >= $3::timestamptz
and ($4::text is null or source = $4)",
)
.bind(query.workspace_id.as_str())
.bind(operation_id.as_str())
.bind(query.created_after)
.bind(
query
.source
.as_ref()
.map(|value| serialize_enum_text(value, "source"))
.transpose()?,
)
.fetch_one(&self.pool)
.await?;
let calls_total = to_u64(row.try_get::<i64, _>("calls_total")?, "calls_total")?;
if calls_total == 0 {
return Ok(None);
}
Ok(Some(UsageRollupRecord {
rollup: UsageRollup {
workspace_id: query.workspace_id.clone(),
agent_id: None,
operation_id: Some(operation_id.clone()),
period: query.period,
calls_total,
calls_ok: to_u64(row.try_get::<i64, _>("calls_ok")?, "calls_ok")?,
calls_error: to_u64(row.try_get::<i64, _>("calls_error")?, "calls_error")?,
p50_ms: to_u64(row.try_get::<i64, _>("p50_ms")?, "p50_ms")?,
p95_ms: to_u64(row.try_get::<i64, _>("p95_ms")?, "p95_ms")?,
p99_ms: to_u64(row.try_get::<i64, _>("p99_ms")?, "p99_ms")?,
},
}))
}
pub async fn list_usage_by_agent(
&self,
query: UsageQuery<'_>,
) -> Result<Vec<UsageAgentBreakdown>, RegistryError> {
let rows = sqlx::query(
"select
a.id as agent_id,
a.slug as agent_slug,
a.display_name as agent_display_name,
count(*)::bigint as calls_total,
count(*) filter (where l.status = 'error')::bigint as calls_error,
coalesce(percentile_cont(0.5) within group (order by l.duration_ms), 0)::bigint as p50_ms,
coalesce(percentile_cont(0.95) within group (order by l.duration_ms), 0)::bigint as p95_ms,
coalesce(percentile_cont(0.99) within group (order by l.duration_ms), 0)::bigint as p99_ms
from invocation_logs l
join agents a on a.id = l.agent_id
where l.workspace_id = $1
and l.created_at >= $2::timestamptz
and ($3::text is null or l.source = $3)
group by a.id, a.slug, a.display_name
order by calls_total desc, a.slug asc",
)
.bind(query.workspace_id.as_str())
.bind(query.created_after)
.bind(
query
.source
.as_ref()
.map(|value| serialize_enum_text(value, "source"))
.transpose()?,
)
.fetch_all(&self.pool)
.await?;
rows.iter().map(map_usage_agent_breakdown).collect()
}
pub async fn get_usage_for_agent(
&self,
query: UsageQuery<'_>,
agent_id: &AgentId,
) -> Result<Option<UsageRollupRecord>, RegistryError> {
let row = sqlx::query(
"select
count(*)::bigint as calls_total,
count(*) filter (where status = 'ok')::bigint as calls_ok,
count(*) filter (where status = 'error')::bigint as calls_error,
coalesce(percentile_cont(0.5) within group (order by duration_ms), 0)::bigint as p50_ms,
coalesce(percentile_cont(0.95) within group (order by duration_ms), 0)::bigint as p95_ms,
coalesce(percentile_cont(0.99) within group (order by duration_ms), 0)::bigint as p99_ms
from invocation_logs
where workspace_id = $1
and agent_id = $2
and created_at >= $3::timestamptz
and ($4::text is null or source = $4)",
)
.bind(query.workspace_id.as_str())
.bind(agent_id.as_str())
.bind(query.created_after)
.bind(
query
.source
.as_ref()
.map(|value| serialize_enum_text(value, "source"))
.transpose()?,
)
.fetch_one(&self.pool)
.await?;
let calls_total = to_u64(row.try_get::<i64, _>("calls_total")?, "calls_total")?;
if calls_total == 0 {
return Ok(None);
}
Ok(Some(UsageRollupRecord {
rollup: UsageRollup {
workspace_id: query.workspace_id.clone(),
agent_id: Some(agent_id.clone()),
operation_id: None,
period: query.period,
calls_total,
calls_ok: to_u64(row.try_get::<i64, _>("calls_ok")?, "calls_ok")?,
calls_error: to_u64(row.try_get::<i64, _>("calls_error")?, "calls_error")?,
p50_ms: to_u64(row.try_get::<i64, _>("p50_ms")?, "p50_ms")?,
p95_ms: to_u64(row.try_get::<i64, _>("p95_ms")?, "p95_ms")?,
p99_ms: to_u64(row.try_get::<i64, _>("p99_ms")?, "p99_ms")?,
},
}))
}
pub async fn get_workspace(
&self,
workspace_id: &WorkspaceId,
@@ -1689,6 +2134,45 @@ fn map_platform_api_key_record(row: &PgRow) -> Result<PlatformApiKeyRecord, Regi
})
}
fn map_invocation_log_record(row: &PgRow) -> Result<InvocationLogRecord, RegistryError> {
Ok(InvocationLogRecord {
log: InvocationLog {
id: InvocationLogId::new(row.try_get::<String, _>("id")?),
workspace_id: WorkspaceId::new(row.try_get::<String, _>("workspace_id")?),
agent_id: row
.try_get::<Option<String>, _>("agent_id")?
.map(AgentId::new),
operation_id: OperationId::new(row.try_get::<String, _>("operation_id")?),
source: deserialize_enum_text(&row.try_get::<String, _>("source")?, "source")?,
level: deserialize_enum_text(&row.try_get::<String, _>("level")?, "level")?,
status: deserialize_enum_text(&row.try_get::<String, _>("status")?, "status")?,
tool_name: row.try_get("tool_name")?,
message: row.try_get("message")?,
request_id: row.try_get("request_id")?,
status_code: match row.try_get::<Option<i32>, _>("status_code")? {
Some(value) => {
Some(
u16::try_from(value).map_err(|_| RegistryError::InvalidNumericValue {
field: "status_code",
value: i64::from(value),
})?,
)
}
None => None,
},
duration_ms: to_u64(row.try_get::<i64, _>("duration_ms")?, "duration_ms")?,
error_kind: row.try_get("error_kind")?,
request_preview: row.try_get::<Json<Value>, _>("request_preview_json")?.0,
response_preview: row.try_get::<Json<Value>, _>("response_preview_json")?.0,
created_at: row.try_get("created_at")?,
},
operation_name: row.try_get("operation_name")?,
operation_display_name: row.try_get("operation_display_name")?,
agent_slug: row.try_get("agent_slug")?,
agent_display_name: row.try_get("agent_display_name")?,
})
}
fn map_agent_summary(row: &PgRow) -> Result<AgentSummary, RegistryError> {
Ok(AgentSummary {
id: AgentId::new(row.try_get::<String, _>("id")?),
@@ -1918,6 +2402,33 @@ fn map_yaml_import_job(row: &PgRow) -> Result<YamlImportJob, RegistryError> {
})
}
fn map_usage_operation_breakdown(row: &PgRow) -> Result<UsageOperationBreakdown, RegistryError> {
Ok(UsageOperationBreakdown {
operation_id: OperationId::new(row.try_get::<String, _>("operation_id")?),
operation_name: row.try_get("operation_name")?,
operation_display_name: row.try_get("operation_display_name")?,
protocol: deserialize_enum_text(&row.try_get::<String, _>("protocol")?, "protocol")?,
calls_total: to_u64(row.try_get::<i64, _>("calls_total")?, "calls_total")?,
calls_error: to_u64(row.try_get::<i64, _>("calls_error")?, "calls_error")?,
p50_ms: to_u64(row.try_get::<i64, _>("p50_ms")?, "p50_ms")?,
p95_ms: to_u64(row.try_get::<i64, _>("p95_ms")?, "p95_ms")?,
p99_ms: to_u64(row.try_get::<i64, _>("p99_ms")?, "p99_ms")?,
})
}
fn map_usage_agent_breakdown(row: &PgRow) -> Result<UsageAgentBreakdown, RegistryError> {
Ok(UsageAgentBreakdown {
agent_id: AgentId::new(row.try_get::<String, _>("agent_id")?),
agent_slug: row.try_get("agent_slug")?,
agent_display_name: row.try_get("agent_display_name")?,
calls_total: to_u64(row.try_get::<i64, _>("calls_total")?, "calls_total")?,
calls_error: to_u64(row.try_get::<i64, _>("calls_error")?, "calls_error")?,
p50_ms: to_u64(row.try_get::<i64, _>("p50_ms")?, "p50_ms")?,
p95_ms: to_u64(row.try_get::<i64, _>("p95_ms")?, "p95_ms")?,
p99_ms: to_u64(row.try_get::<i64, _>("p99_ms")?, "p99_ms")?,
})
}
fn serialize_json_value<T: Serialize>(value: &T) -> Result<Value, RegistryError> {
Ok(serde_json::to_value(value)?)
}
@@ -1961,6 +2472,19 @@ fn from_db_version(value: i32, field: &'static str) -> Result<u32, RegistryError
})
}
fn to_u64(value: i64, field: &'static str) -> Result<u64, RegistryError> {
u64::try_from(value).map_err(|_| RegistryError::InvalidNumericValue { field, value })
}
fn usage_bucket_sql(bucket: crate::model::UsageBucket) -> &'static str {
match bucket {
crate::model::UsageBucket::Hour => "hour",
crate::model::UsageBucket::Day => "day",
crate::model::UsageBucket::Week => "week",
crate::model::UsageBucket::Month => "month",
}
}
#[cfg(test)]
mod tests {
use std::{
+16 -1
View File
@@ -35,11 +35,26 @@ impl RuntimeExecutor {
operation: &RuntimeOperation,
input: &Value,
) -> Result<Value, RuntimeError> {
let prepared_request = self.prepare_request(operation, input)?;
self.execute_prepared(operation, prepared_request).await
}
pub fn prepare_request(
&self,
operation: &RuntimeOperation,
input: &Value,
) -> Result<PreparedRequest, RuntimeError> {
operation.input_schema.validate_shape(input)?;
let mapped_input = operation.input_mapping.apply(&json!({ "mcp": input }))?;
let prepared_request = PreparedRequest::from_mapping_output(&mapped_input)?;
PreparedRequest::from_mapping_output(&mapped_input)
}
async fn execute_prepared(
&self,
operation: &RuntimeOperation,
prepared_request: PreparedRequest,
) -> Result<Value, RuntimeError> {
let adapter_response = match &operation.target {
Target::Grpc(target) => {
let request = GrpcRequest {