1use 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 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}