Skip to main content

skiff_cli/mcp/
http.rs

1//! MCP streamable HTTP client via rmcp.
2
3use std::sync::Arc;
4use std::time::Duration;
5
6use rmcp::{
7    transport::streamable_http_client::StreamableHttpClientTransportConfig,
8    transport::StreamableHttpClientTransport, ServiceExt,
9};
10use serde_json::{Map, Value};
11
12use crate::cache::load_cached;
13use crate::error::{Error, Result};
14use crate::mcp::common::{auth_headers_to_http, call_tool_on, list_tools_on, McpClient};
15use crate::oauth::OAuthReady;
16use crate::tools_index::save_tools_and_index;
17
18pub async fn fetch_mcp_tools_http(
19    url: &str,
20    auth_headers: &[(String, String)],
21    cache_key: &str,
22    ttl: u64,
23    refresh: bool,
24    oauth: Option<&OAuthReady>,
25) -> Result<Vec<Value>> {
26    let tools_key = format!("{cache_key}_tools");
27    if !refresh {
28        if let Some(cached) = load_cached(&tools_key, ttl)? {
29            if let Some(arr) = cached.as_array() {
30                let _ = crate::tools_index::save_index(&tools_key, arr);
31                return Ok(arr.clone());
32            }
33        }
34    }
35
36    let tools = list_tools_http(url, auth_headers, oauth).await?;
37    save_tools_and_index(&tools_key, &tools)?;
38    Ok(tools)
39}
40
41pub async fn list_tools_http(
42    url: &str,
43    auth_headers: &[(String, String)],
44    oauth: Option<&OAuthReady>,
45) -> Result<Vec<Value>> {
46    let client = connect_streamable(url, auth_headers, oauth).await?;
47    let tools = list_tools_on(&client).await?;
48    let _ = client.cancel().await;
49    Ok(tools)
50}
51
52pub async fn call_tool_http(
53    url: &str,
54    auth_headers: &[(String, String)],
55    tool_name: &str,
56    arguments: Map<String, Value>,
57    full_envelope: bool,
58    oauth: Option<&OAuthReady>,
59) -> Result<Value> {
60    let client = connect_streamable(url, auth_headers, oauth).await?;
61    let result = call_tool_on(&client, tool_name, arguments, full_envelope).await?;
62    let _ = client.cancel().await;
63    Ok(result)
64}
65
66pub async fn connect_streamable(
67    url: &str,
68    auth_headers: &[(String, String)],
69    oauth: Option<&OAuthReady>,
70) -> Result<McpClient> {
71    let mut headers = auth_headers.to_vec();
72    if let Some(oauth) = oauth {
73        // Refresh-aware access token; inject as Bearer for this one-shot connection.
74        let token = oauth.access_token().await?;
75        headers.retain(|(k, _)| !k.eq_ignore_ascii_case("authorization"));
76        headers.push(("Authorization".into(), format!("Bearer {token}")));
77    }
78    let custom_headers = auth_headers_to_http(&headers)?;
79    let config = StreamableHttpClientTransportConfig::with_uri(Arc::<str>::from(url))
80        .custom_headers(custom_headers);
81
82    let transport = StreamableHttpClientTransport::from_config(config);
83    tokio::time::timeout(Duration::from_secs(30), ().serve(transport))
84        .await
85        .map_err(|_| Error::runtime("MCP HTTP initialize timed out after 30s"))?
86        .map_err(|e| Error::runtime(format!("MCP HTTP initialize failed: {e}")))
87}