systemprompt-runtime 0.64.0

Application runtime for systemprompt.io AI governance infrastructure. AppContext, lifecycle builder, extension registry, and module wiring for the MCP governance pipeline.
Documentation
//! MCP tool-call trace queries.
//!
//! Copyright (c) systemprompt.io — Business Source License 1.1.
//! See <https://systemprompt.io> for licensing details.

use systemprompt_identifiers::{
    AiRequestId, ArtifactId, ContextId, McpExecutionId, McpServerId, McpToolName, TaskId,
};

use super::{Result, TraceRepository};
use crate::trace::models::{AiRequestInfo, McpToolExecution, TaskArtifact, ToolLogEntry};

impl TraceRepository {
    pub async fn fetch_mcp_executions(
        &self,
        task_id: &TaskId,
        context_id: &ContextId,
    ) -> Result<Vec<McpToolExecution>> {
        let rows = sqlx::query!(
            r#"SELECT mcp_execution_id, tool_name, server_name, status, execution_time_ms,
                  error_message, input, output
           FROM mcp_tool_executions
           WHERE task_id = $1 OR context_id = $2
           ORDER BY started_at"#,
            task_id.as_str(),
            context_id.as_str()
        )
        .fetch_all(&*self.pool)
        .await?;

        Ok(rows
            .into_iter()
            .map(|r| McpToolExecution {
                mcp_execution_id: McpExecutionId::new(r.mcp_execution_id),
                tool_name: McpToolName::new(r.tool_name),
                server_name: McpServerId::new(r.server_name),
                status: r.status,
                execution_time_ms: r.execution_time_ms,
                error_message: r.error_message,
                input: r.input,
                output: r.output,
            })
            .collect())
    }

    pub async fn fetch_mcp_linked_ai_requests(
        &self,
        mcp_execution_id: &McpExecutionId,
    ) -> Result<Vec<AiRequestInfo>> {
        let rows = sqlx::query!(
        r#"SELECT id, model, provider, max_tokens, input_tokens, output_tokens, cost_microdollars, latency_ms
           FROM ai_requests
           WHERE mcp_execution_id = $1
           ORDER BY created_at"#,
        mcp_execution_id.as_str()
    )
    .fetch_all(&*self.pool)
    .await?;

        Ok(rows
            .into_iter()
            .map(|r| AiRequestInfo {
                id: AiRequestId::new(r.id),
                provider: r.provider,
                model: r.model,
                max_tokens: r.max_tokens,
                input_tokens: r.input_tokens,
                output_tokens: r.output_tokens,
                cost_microdollars: r.cost_microdollars,
                latency_ms: r.latency_ms,
            })
            .collect())
    }

    pub async fn fetch_tool_logs(
        &self,
        task_id: &TaskId,
        context_id: &ContextId,
    ) -> Result<Vec<ToolLogEntry>> {
        let rows = sqlx::query!(
        r#"SELECT timestamp, level, module, message
           FROM logs
           WHERE (task_id = $1 OR context_id = $2)
             AND (
                 (module LIKE '%_tools' OR module LIKE '%_manager' OR module LIKE 'create_%' OR module LIKE 'update_%' OR module LIKE 'research_%')
                 OR (level = 'ERROR' AND message LIKE '%tool%')
                 OR message LIKE 'Tool executed%'
                 OR message LIKE 'Tool failed%'
                 OR message LIKE 'MCP execution%'
             )
           ORDER BY timestamp"#,
        task_id.as_str(),
        context_id.as_str()
    )
    .fetch_all(&*self.pool)
    .await?;

        rows.into_iter()
            .map(|r| -> Result<ToolLogEntry> {
                Ok(ToolLogEntry {
                    timestamp: r.timestamp,
                    level: r.level.parse()?,
                    module: r.module,
                    message: r.message,
                })
            })
            .collect()
    }

    pub async fn fetch_task_artifacts(
        &self,
        task_id: &TaskId,
        context_id: &ContextId,
    ) -> Result<Vec<TaskArtifact>> {
        let rows = sqlx::query!(
        r#"SELECT ta.artifact_id, ta.artifact_type, ta.name, ta.source, ta.tool_name,
                  ap.part_kind as "part_kind?", ap.text_content as "text_content?",
                  ap.data_content as "data_content?"
           FROM task_artifacts ta
           LEFT JOIN artifact_parts ap ON ta.artifact_id = ap.artifact_id AND ta.context_id = ap.context_id
           WHERE ta.task_id = $1 OR ta.context_id = $2
           ORDER BY ta.created_at, ap.sequence_number"#,
        task_id.as_str(),
        context_id.as_str()
    )
    .fetch_all(&*self.pool)
    .await?;

        Ok(rows
            .into_iter()
            .map(|r| TaskArtifact {
                artifact_id: ArtifactId::new(r.artifact_id),
                artifact_type: r.artifact_type,
                name: r.name,
                source: r.source,
                tool_name: r.tool_name.map(McpToolName::new),
                part_kind: r.part_kind,
                text_content: r.text_content,
                data_content: r.data_content,
            })
            .collect())
    }
}