use super::*;
use crate::value::VmDictExt;
pub(crate) fn vm_value_to_serde(val: &VmValue) -> serde_json::Value {
match val {
VmValue::String(s) => serde_json::Value::String(s.to_string()),
VmValue::Int(n) => serde_json::json!(*n),
VmValue::Float(n) => serde_json::json!(*n),
VmValue::Decimal(d) => serde_json::Value::String(d.to_string()),
VmValue::Bool(b) => serde_json::Value::Bool(*b),
VmValue::Nil => serde_json::Value::Null,
VmValue::List(items) => {
serde_json::Value::Array(items.iter().map(vm_value_to_serde).collect())
}
VmValue::Dict(map) => {
let obj: serde_json::Map<String, serde_json::Value> = map
.iter()
.map(|(k, v)| (k.to_string(), vm_value_to_serde(v)))
.collect();
serde_json::Value::Object(obj)
}
_ => serde_json::Value::Null,
}
}
pub(crate) async fn call_mcp_tool(
client: &VmMcpClientHandle,
tool_name: &str,
arguments: serde_json::Value,
) -> Result<serde_json::Value, VmError> {
let (content, _hint) = call_mcp_tool_with_hint(client, tool_name, arguments).await?;
Ok(content)
}
pub(crate) async fn call_mcp_tool_with_hint(
client: &VmMcpClientHandle,
tool_name: &str,
arguments: serde_json::Value,
) -> Result<(serde_json::Value, Option<McpCacheHint>), VmError> {
let progress_token = client_progress::issue_token(&client.name, tool_name);
let _progress_guard = progress_token
.as_ref()
.map(|tok| client_progress::ProgressTokenGuard {
token: tok.token.clone(),
});
let mut result = client
.call(
"tools/call",
tool_call_params(tool_name, arguments, progress_token.as_ref()),
)
.await?;
if result.get("resultType").and_then(|value| value.as_str()) == Some("task") {
result = resolve_task_result(client, result).await?;
}
if result.get("isError").and_then(|v| v.as_bool()) == Some(true) {
let error_text = extract_content_text(&result);
return Err(VmError::Thrown(VmValue::String(arcstr::ArcStr::from(
error_text,
))));
}
let cache_hint = McpCacheHint::from_result(&result);
if let Some(structured) = result.get("structuredContent") {
return Ok((structured.clone(), cache_hint));
}
let content = result
.get("content")
.and_then(|c| c.as_array())
.cloned()
.unwrap_or_default();
let unwrapped =
if content.len() == 1 && content[0].get("type").and_then(|t| t.as_str()) == Some("text") {
if let Some(text) = content[0].get("text").and_then(|t| t.as_str()) {
serde_json::Value::String(text.to_string())
} else if content.is_empty() {
serde_json::Value::Null
} else {
serde_json::Value::Array(content)
}
} else if content.is_empty() {
serde_json::Value::Null
} else {
serde_json::Value::Array(content)
};
Ok((unwrapped, cache_hint))
}
async fn resolve_task_result(
client: &VmMcpClientHandle,
created: serde_json::Value,
) -> Result<serde_json::Value, VmError> {
let task_id = created
.get("taskId")
.and_then(serde_json::Value::as_str)
.ok_or_else(|| VmError::Runtime("MCP task result is missing taskId".into()))?
.to_string();
let poll_interval_ms = created
.get("pollIntervalMs")
.and_then(serde_json::Value::as_u64)
.unwrap_or(crate::mcp_protocol::DEFAULT_TASK_POLL_INTERVAL_MS)
.clamp(10, 5_000);
tokio::time::timeout(MCP_TIMEOUT, async {
loop {
let task = client
.call(
crate::mcp_protocol::METHOD_TASKS_GET,
serde_json::json!({"taskId": task_id}),
)
.await?;
match task.get("status").and_then(serde_json::Value::as_str) {
Some("completed") => {
return task.get("result").cloned().ok_or_else(|| {
VmError::Runtime(format!("MCP task '{task_id}' completed without a result"))
});
}
Some("failed") => {
let error = task.get("error").cloned().unwrap_or_else(|| {
serde_json::json!({
"code": -32603,
"message": format!("MCP task '{task_id}' failed"),
})
});
return Err(jsonrpc_error_to_vm_error(&error));
}
Some("cancelled") => {
return Err(VmError::Runtime(format!(
"MCP task '{task_id}' was cancelled"
)));
}
Some("input_required") => {
let Some(input_round) = resolve_input_required_result(client, &task).await?
else {
return Err(VmError::Runtime(format!(
"MCP task '{task_id}' requires input but supplied no inputRequests"
)));
};
client
.call(
crate::mcp_protocol::METHOD_TASKS_UPDATE,
serde_json::json!({
"taskId": task_id,
"inputResponses": input_round.input_responses,
}),
)
.await?;
}
Some("working") => {}
Some(status) => {
return Err(VmError::Runtime(format!(
"MCP task '{task_id}' returned unknown status {status:?}"
)));
}
None => {
return Err(VmError::Runtime(format!(
"MCP task '{task_id}' response is missing status"
)));
}
}
tokio::time::sleep(std::time::Duration::from_millis(poll_interval_ms)).await;
}
})
.await
.map_err(|_| {
VmError::Runtime(format!(
"MCP task '{task_id}' did not complete within {}s",
MCP_TIMEOUT.as_secs()
))
})?
}
pub(crate) fn tool_call_params(
tool_name: &str,
arguments: serde_json::Value,
progress_token: Option<&client_progress::ProgressTokenContext>,
) -> serde_json::Value {
let mut params = serde_json::Map::from_iter([
(
"name".to_string(),
serde_json::Value::String(tool_name.to_string()),
),
("arguments".to_string(), arguments),
]);
if let Some(tok) = progress_token {
params.insert(
"_meta".to_string(),
serde_json::json!({ "progressToken": tok.token }),
);
}
serde_json::Value::Object(params)
}
pub(crate) fn with_input_round(
params: serde_json::Value,
input_round: McpInputRound,
) -> Result<serde_json::Value, VmError> {
let mut params = params.as_object().cloned().ok_or_else(|| {
VmError::Runtime("MCP MRTR retry parameters must be an object".to_string())
})?;
params.insert("inputResponses".to_string(), input_round.input_responses);
if let Some(request_state) = input_round.request_state {
params.insert("requestState".to_string(), request_state);
}
Ok(serde_json::Value::Object(params))
}
pub(crate) async fn resolve_input_required_result(
client: &VmMcpClientHandle,
result: &serde_json::Value,
) -> Result<Option<McpInputRound>, VmError> {
let Some(input_requests) = result.get("inputRequests") else {
if result.get("requestState").is_some() {
return Ok(Some(McpInputRound {
input_responses: serde_json::json!({}),
request_state: result.get("requestState").cloned(),
}));
}
return Ok(None);
};
let input_requests: rmcp::model::InputRequests = serde_json::from_value(input_requests.clone())
.map_err(|error| VmError::Runtime(format!("invalid MCP inputRequests: {error}")))?;
let fixtures = client.capability_fixtures().await;
let mut responses = rmcp::model::InputResponses::new();
for (key, input_request) in input_requests {
let mut request = serde_json::to_value(input_request).map_err(|error| {
VmError::Runtime(format!("invalid MCP input request {key:?}: {error}"))
})?;
let object = request.as_object_mut().ok_or_else(|| {
VmError::Runtime(format!("MCP input request {key:?} is not an object"))
})?;
object.insert("jsonrpc".to_string(), serde_json::json!("2.0"));
object.insert("id".to_string(), serde_json::json!(format!("input-{key}")));
let response = resolve_embedded_input_request(&client.name, &request, fixtures.as_deref())
.await
.ok_or_else(|| {
VmError::Runtime(format!(
"MCP input_required request {key:?} used an unsupported method"
))
})?;
if let Some(result) = response.get("result") {
responses.insert(key.clone(), result.clone());
} else if let Some(error) = response.get("error") {
return Err(jsonrpc_error_to_vm_error(error));
} else {
responses.insert(key.clone(), serde_json::Value::Null);
}
}
Ok(Some(McpInputRound {
input_responses: serde_json::to_value(responses)
.expect("MCP input responses are JSON serializable"),
request_state: result.get("requestState").cloned(),
}))
}
pub fn register_mcp_builtins(vm: &mut Vm) {
vm.register_capability_method(
harn_builtin_meta::CapabilityId::Tools,
"mcp_roots",
mcp_roots_builtin,
);
register_supervised_mcp_host_builtins(vm);
vm.register_async_capability_method(
harn_builtin_meta::CapabilityId::Tools,
"mcp_connect",
|_ctx, args| async move {
let command = args.first().map(|a| a.display()).unwrap_or_default();
if command.is_empty() {
return Err(VmError::Thrown(VmValue::String(arcstr::ArcStr::from(
"mcp_connect: command is required",
))));
}
let cmd_args: Vec<String> = match args.get(1) {
Some(VmValue::List(list)) => list.iter().map(|v| v.display()).collect(),
_ => Vec::new(),
};
let options = mcp_connect_options(args.get(2))?;
let handle = mcp_connect_stdio_impl(
&command,
&cmd_args,
&BTreeMap::new(),
options.protocol_version,
)
.await?;
Ok(VmValue::mcp_client(handle))
},
);
vm.register_async_capability_method(
harn_builtin_meta::CapabilityId::Tools,
"mcp_ensure_active",
|_ctx, args| async move {
let name = match args.first() {
Some(VmValue::String(s)) => s.to_string(),
Some(other) => other.display(),
None => String::new(),
};
if name.is_empty() {
return Err(VmError::Thrown(VmValue::String(arcstr::ArcStr::from(
"mcp_ensure_active: server name is required",
))));
}
let handle = crate::mcp_registry::ensure_active(&name).await?;
Ok(VmValue::mcp_client(handle))
},
);
vm.register_capability_method(
harn_builtin_meta::CapabilityId::Tools,
"mcp_release",
|args, _out| {
let name = match args.first() {
Some(VmValue::String(s)) => s.to_string(),
Some(other) => other.display(),
None => {
return Err(VmError::Thrown(VmValue::String(arcstr::ArcStr::from(
"mcp_release: server name is required",
))));
}
};
crate::mcp_registry::release(&name);
Ok(VmValue::Nil)
},
);
vm.register_capability_method(
harn_builtin_meta::CapabilityId::Tools,
"mcp_registry_status",
|_args, _out| {
let mut out = Vec::new();
for entry in crate::mcp_registry::snapshot_status() {
let mut dict = BTreeMap::new();
dict.put_str("name", entry.name.as_str());
dict.put_str("transport", entry.transport.as_str());
if let Some(url) = entry.url {
dict.put_str("url", url.as_str());
} else {
dict.insert("url".to_string(), VmValue::Nil);
}
dict.insert("lazy".to_string(), VmValue::Bool(entry.lazy));
dict.insert("active".to_string(), VmValue::Bool(entry.active));
dict.insert(
"ref_count".to_string(),
VmValue::Int(entry.ref_count as i64),
);
if let Some(card) = entry.card {
dict.put_str("card", card.as_str());
}
out.push(VmValue::dict(dict));
}
Ok(VmValue::List(std::sync::Arc::new(out)))
},
);
vm.register_async_capability_method(
harn_builtin_meta::CapabilityId::Tools,
"mcp_reauth_expired",
|_ctx, _args| async move { reauth_expired_impl().await },
);
vm.register_async_capability_method(
harn_builtin_meta::CapabilityId::Tools,
"mcp_server_card",
|_ctx, args| async move {
let target = match args.first() {
Some(VmValue::String(s)) => s.to_string(),
Some(other) => other.display(),
None => {
return Err(VmError::Thrown(VmValue::String(arcstr::ArcStr::from(
"mcp_server_card: server name, URL, or path is required",
))));
}
};
let source = if target.starts_with("http://")
|| target.starts_with("https://")
|| target.contains('/')
|| target.contains('\\')
|| target.ends_with(".json")
{
target.clone()
} else {
match crate::mcp_registry::get_registration(&target) {
Some(reg) => match reg.card {
Some(card) => card,
None => {
return Err(VmError::Thrown(VmValue::String(arcstr::ArcStr::from(
format!(
"mcp_server_card: server '{target}' has no 'card' field in harn.toml"
),
))));
}
},
None => {
return Err(VmError::Thrown(VmValue::String(arcstr::ArcStr::from(
format!(
"mcp_server_card: no MCP server '{target}' registered (check harn.toml) \
— pass a URL or path directly instead"
),
))));
}
}
};
let card = crate::mcp_card::fetch_server_card(&source, None)
.await
.map_err(|e| {
VmError::Thrown(VmValue::String(arcstr::ArcStr::from(format!(
"mcp_server_card: {e}"
))))
})?;
Ok(json_to_vm_value(&card))
},
);
vm.register_async_capability_method(
harn_builtin_meta::CapabilityId::Tools,
"mcp_list_tools",
|_ctx, args| async move {
let client = match args.first() {
Some(VmValue::McpClient(c)) => c.clone(),
_ => {
return Err(VmError::Thrown(VmValue::String(arcstr::ArcStr::from(
"mcp_list_tools: argument must be an MCP client",
))));
}
};
let result = client.call("tools/list", serde_json::json!({})).await?;
client.record_cache_hint("tools/list", &result).await;
let mut tools = result
.get("tools")
.and_then(|t| t.as_array())
.cloned()
.unwrap_or_default();
tools = filter_tools_for_client(&tools);
client.store_http_tool_headers(&tools).await;
let server_name = client.name.clone();
for tool in tools.iter_mut() {
if let Some(obj) = tool.as_object_mut() {
obj.entry("_mcp_server")
.or_insert_with(|| serde_json::Value::String(server_name.clone()));
}
}
let vm_tools: Vec<VmValue> = tools.iter().map(json_to_vm_value).collect();
Ok(VmValue::List(std::sync::Arc::new(vm_tools)))
},
);
vm.register_async_capability_method(
harn_builtin_meta::CapabilityId::Tools,
"mcp_call",
|_ctx, args| async move {
let client = match args.first() {
Some(VmValue::McpClient(c)) => c.clone(),
_ => {
return Err(VmError::Thrown(VmValue::String(arcstr::ArcStr::from(
"mcp_call: first argument must be an MCP client",
))));
}
};
let tool_name = args.get(1).map(|a| a.display()).unwrap_or_default();
if tool_name.is_empty() {
return Err(VmError::Thrown(VmValue::String(arcstr::ArcStr::from(
"mcp_call: tool name is required",
))));
}
let arguments = match args.get(2) {
Some(VmValue::Dict(d)) => {
let obj: serde_json::Map<String, serde_json::Value> = d
.iter()
.map(|(k, v)| (k.to_string(), vm_value_to_serde(v)))
.collect();
serde_json::Value::Object(obj)
}
_ => serde_json::json!({}),
};
Ok(json_to_vm_value(
&call_mcp_tool(&client, &tool_name, arguments).await?,
))
},
);
vm.register_async_capability_method(
harn_builtin_meta::CapabilityId::Tools,
"mcp_server_info",
|_ctx, args| async move {
let client = match args.first() {
Some(VmValue::McpClient(c)) => c.clone(),
_ => {
return Err(VmError::Thrown(VmValue::String(arcstr::ArcStr::from(
"mcp_server_info: argument must be an MCP client",
))));
}
};
let guard = client.inner.lock().await;
if guard.is_none() {
return Err(VmError::Runtime("MCP client is disconnected".into()));
}
drop(guard);
let mut info = BTreeMap::new();
info.put_str("name", client.name.as_str());
info.insert("connected".to_string(), VmValue::Bool(true));
let discovery = client
.discovery_result
.lock()
.await
.clone()
.unwrap_or(serde_json::Value::Null);
if !discovery.is_null() {
if let Some(instructions) = discovery
.get("instructions")
.or_else(|| {
discovery
.get("serverInfo")
.and_then(|value| value.get("instructions"))
})
.and_then(|value| value.as_str())
.filter(|value| !value.is_empty())
{
info.put_str("instructions", instructions);
}
info.insert("discovery".to_string(), json_to_vm_value(&discovery));
}
let cache_hints = client.cache_hints.lock().await;
if !cache_hints.is_empty() {
info.insert(
"cache_hints".to_string(),
json_to_vm_value(&cache_hints_to_json(cache_hints.iter())),
);
}
Ok(VmValue::dict(info))
},
);
vm.register_async_capability_method(
harn_builtin_meta::CapabilityId::Tools,
"mcp_disconnect",
|_ctx, args| async move {
let client = match args.first() {
Some(VmValue::McpClient(c)) => c.clone(),
_ => {
return Err(VmError::Thrown(VmValue::String(arcstr::ArcStr::from(
"mcp_disconnect: argument must be an MCP client",
))));
}
};
client.disconnect().await?;
Ok(VmValue::Nil)
},
);
vm.register_async_capability_method(
harn_builtin_meta::CapabilityId::Tools,
"mcp_list_resources",
|_ctx, args| async move {
let client = match args.first() {
Some(VmValue::McpClient(c)) => c.clone(),
_ => {
return Err(VmError::Thrown(VmValue::String(arcstr::ArcStr::from(
"mcp_list_resources: argument must be an MCP client",
))));
}
};
let result = client.call("resources/list", serde_json::json!({})).await?;
client.record_cache_hint("resources/list", &result).await;
let resources = result
.get("resources")
.and_then(|r| r.as_array())
.cloned()
.unwrap_or_default();
let vm_resources: Vec<VmValue> = resources.iter().map(json_to_vm_value).collect();
Ok(VmValue::List(std::sync::Arc::new(vm_resources)))
},
);
vm.register_async_capability_method(
harn_builtin_meta::CapabilityId::Tools,
"mcp_read_resource",
|_ctx, args| async move {
let client = match args.first() {
Some(VmValue::McpClient(c)) => c.clone(),
_ => {
return Err(VmError::Thrown(VmValue::String(arcstr::ArcStr::from(
"mcp_read_resource: first argument must be an MCP client",
))));
}
};
let uri = args.get(1).map(|a| a.display()).unwrap_or_default();
if uri.is_empty() {
return Err(VmError::Thrown(VmValue::String(arcstr::ArcStr::from(
"mcp_read_resource: URI is required",
))));
}
let result = client
.call("resources/read", serde_json::json!({ "uri": uri }))
.await?;
client.record_cache_hint("resources/read", &result).await;
let contents = result
.get("contents")
.and_then(|c| c.as_array())
.cloned()
.unwrap_or_default();
if contents.len() == 1 {
if let Some(text) = contents[0].get("text").and_then(|t| t.as_str()) {
return Ok(VmValue::String(arcstr::ArcStr::from(text)));
}
}
if contents.is_empty() {
Ok(VmValue::Nil)
} else {
Ok(VmValue::List(std::sync::Arc::new(
contents.iter().map(json_to_vm_value).collect(),
)))
}
},
);
vm.register_async_capability_method(
harn_builtin_meta::CapabilityId::Tools,
"mcp_list_resource_templates",
|_ctx, args| async move {
let client = match args.first() {
Some(VmValue::McpClient(c)) => c.clone(),
_ => {
return Err(VmError::Thrown(VmValue::String(arcstr::ArcStr::from(
"mcp_list_resource_templates: argument must be an MCP client",
))));
}
};
let result = client
.call("resources/templates/list", serde_json::json!({}))
.await?;
client
.record_cache_hint("resources/templates/list", &result)
.await;
let templates = result
.get("resourceTemplates")
.and_then(|r| r.as_array())
.cloned()
.unwrap_or_default();
let vm_templates: Vec<VmValue> = templates.iter().map(json_to_vm_value).collect();
Ok(VmValue::List(std::sync::Arc::new(vm_templates)))
},
);
vm.register_async_capability_method(
harn_builtin_meta::CapabilityId::Tools,
"mcp_list_prompts",
|_ctx, args| async move {
let client = match args.first() {
Some(VmValue::McpClient(c)) => c.clone(),
_ => {
return Err(VmError::Thrown(VmValue::String(arcstr::ArcStr::from(
"mcp_list_prompts: argument must be an MCP client",
))));
}
};
let result = client.call("prompts/list", serde_json::json!({})).await?;
client.record_cache_hint("prompts/list", &result).await;
let prompts = result
.get("prompts")
.and_then(|p| p.as_array())
.cloned()
.unwrap_or_default();
let vm_prompts: Vec<VmValue> = prompts.iter().map(json_to_vm_value).collect();
Ok(VmValue::List(std::sync::Arc::new(vm_prompts)))
},
);
vm.register_async_capability_method(
harn_builtin_meta::CapabilityId::Tools,
"mcp_get_prompt",
|_ctx, args| async move {
let client = match args.first() {
Some(VmValue::McpClient(c)) => c.clone(),
_ => {
return Err(VmError::Thrown(VmValue::String(arcstr::ArcStr::from(
"mcp_get_prompt: first argument must be an MCP client",
))));
}
};
let name = args.get(1).map(|a| a.display()).unwrap_or_default();
if name.is_empty() {
return Err(VmError::Thrown(VmValue::String(arcstr::ArcStr::from(
"mcp_get_prompt: prompt name is required",
))));
}
let arguments = match args.get(2) {
Some(VmValue::Dict(d)) => {
let obj: serde_json::Map<String, serde_json::Value> = d
.iter()
.map(|(k, v)| (k.to_string(), vm_value_to_serde(v)))
.collect();
serde_json::Value::Object(obj)
}
_ => serde_json::json!({}),
};
let result = client
.call(
"prompts/get",
serde_json::json!({
"name": name,
"arguments": arguments,
}),
)
.await?;
Ok(json_to_vm_value(&result))
},
);
}
pub(crate) fn mcp_roots_builtin(_args: &[VmValue], _out: &mut String) -> Result<VmValue, VmError> {
Ok(VmValue::List(std::sync::Arc::new(
current_mcp_roots()
.iter()
.map(|root| json_to_vm_value(&root.script_json()))
.collect(),
)))
}
pub(crate) fn register_supervised_mcp_host_builtins(vm: &mut Vm) {
vm.register_async_capability_method(
harn_builtin_meta::CapabilityId::Tools,
"mcp_host_spawn",
|_ctx, args| async move {
let spec_value = args
.first()
.ok_or_else(|| VmError::Runtime("harn.mcp.spawn: spec dict is required".into()))?;
let spec = vm_value_to_serde(spec_value);
let options: crate::mcp_host::SpawnOptions = match args.get(1) {
Some(VmValue::Dict(_)) => {
let options_json = vm_value_to_serde(args.get(1).unwrap());
serde_json::from_value(options_json).map_err(|e| {
VmError::Runtime(format!("harn.mcp.spawn: invalid options dict: {e}"))
})?
}
_ => Default::default(),
};
let name = crate::mcp_host::spawn(spec, options).await?;
Ok(VmValue::String(arcstr::ArcStr::from(name)))
},
);
vm.register_async_capability_method(
harn_builtin_meta::CapabilityId::Tools,
"mcp_host_tools",
|_ctx, args| async move {
let name = match args.first() {
Some(VmValue::String(s)) => s.to_string(),
Some(other) => other.display(),
None => {
return Err(VmError::Runtime(
"harn.mcp.tools: server name (string) is required".into(),
));
}
};
let tools = crate::mcp_host::tools(&name).await?;
Ok(VmValue::List(std::sync::Arc::new(
tools.iter().map(json_to_vm_value).collect(),
)))
},
);
vm.register_async_capability_method(
harn_builtin_meta::CapabilityId::Tools,
"mcp_host_call",
|_ctx, args| async move {
let name = match args.first() {
Some(VmValue::String(s)) => s.to_string(),
Some(other) => other.display(),
None => {
return Err(VmError::Runtime(
"harn.mcp.call: server name (string) is required".into(),
));
}
};
let tool = match args.get(1) {
Some(VmValue::String(s)) => s.to_string(),
Some(other) => other.display(),
None => {
return Err(VmError::Runtime(
"harn.mcp.call: tool name (string) is required".into(),
));
}
};
let tool_args = match args.get(2) {
Some(VmValue::Dict(d)) => {
let obj: serde_json::Map<String, serde_json::Value> = d
.iter()
.map(|(k, v)| (k.to_string(), vm_value_to_serde(v)))
.collect();
serde_json::Value::Object(obj)
}
Some(VmValue::Nil) | None => serde_json::json!({}),
Some(other) => vm_value_to_serde(other),
};
let result = crate::mcp_host::call(&name, &tool, tool_args).await?;
Ok(json_to_vm_value(&result))
},
);
vm.register_capability_method(
harn_builtin_meta::CapabilityId::Tools,
"mcp_host_stop",
|args, _out| {
let name = match args.first() {
Some(VmValue::String(s)) => s.to_string(),
Some(other) => other.display(),
None => {
return Err(VmError::Runtime(
"harn.mcp.stop: server name (string) is required".into(),
));
}
};
crate::mcp_host::stop(&name)?;
Ok(VmValue::Nil)
},
);
vm.register_capability_method(
harn_builtin_meta::CapabilityId::Tools,
"mcp_host_reload",
|args, _out| {
let name = match args.first() {
Some(VmValue::String(s)) => s.to_string(),
Some(other) => other.display(),
None => {
return Err(VmError::Runtime(
"harn.mcp.reload: server name (string) is required".into(),
));
}
};
crate::mcp_host::reload(&name)?;
Ok(VmValue::Nil)
},
);
vm.register_async_capability_method(
harn_builtin_meta::CapabilityId::Tools,
"mcp_host_discover",
|_ctx, _args| async move {
let entries = crate::mcp_host::discover().await?;
Ok(VmValue::List(std::sync::Arc::new(
entries.iter().map(json_to_vm_value).collect(),
)))
},
);
vm.register_async_capability_method(
harn_builtin_meta::CapabilityId::Tools,
"mcp_host_status",
|_ctx, _args| async move {
let entries = crate::mcp_host::status().await;
let list: Vec<VmValue> = entries
.into_iter()
.map(|s| {
let mut dict = BTreeMap::new();
dict.put_str("name", s.name);
dict.put_str("transport", s.transport);
if let Some(url) = s.url {
dict.put_str("url", url);
} else {
dict.insert("url".to_string(), VmValue::Nil);
}
dict.insert("active".to_string(), VmValue::Bool(s.active));
dict.insert("lazy".to_string(), VmValue::Bool(s.lazy));
dict.insert("ref_count".to_string(), VmValue::Int(s.ref_count as i64));
dict.insert(
"restart_count".to_string(),
VmValue::Int(s.restart_count as i64),
);
dict.insert(
"consecutive_failures".to_string(),
VmValue::Int(s.consecutive_failures as i64),
);
dict.put_str("circuit", s.circuit.as_str());
dict.insert("ejected".to_string(), VmValue::Bool(s.ejected));
dict.insert(
"cache_entries".to_string(),
VmValue::Int(s.cache_entries as i64),
);
if let Some(display_identity) = s.display_identity {
dict.put_str("display_identity", display_identity);
} else {
dict.insert("display_identity".to_string(), VmValue::Nil);
}
VmValue::dict(dict)
})
.collect();
Ok(VmValue::List(std::sync::Arc::new(list)))
},
);
}
const REAUTH_REDIRECT_URI: &str = "http://127.0.0.1:9783/oauth/callback";
fn oauth_servers_from_specs(
specs: &[serde_json::Value],
) -> Vec<crate::mcp_bulk_auth::BulkAuthServer> {
specs
.iter()
.filter_map(|spec| {
let transport = crate::mcp_registry::spec_transport(spec);
let url = spec
.get("url")
.and_then(|value| value.as_str())
.unwrap_or("");
let has_static_bearer = spec
.get("auth_token")
.and_then(|value| value.as_str())
.is_some_and(|token| !token.is_empty());
if transport.as_str() != "http" || url.is_empty() || has_static_bearer {
return None;
}
let name = spec
.get("name")
.and_then(|value| value.as_str())
.unwrap_or(url)
.to_string();
Some(crate::mcp_bulk_auth::BulkAuthServer {
name,
server_url: url.to_string(),
..Default::default()
})
})
.collect()
}
async fn reauth_expired_impl() -> Result<VmValue, VmError> {
let specs: Vec<serde_json::Value> = crate::mcp_registry::all_registrations()
.into_iter()
.map(|registration| registration.spec)
.collect();
let servers = oauth_servers_from_specs(&specs);
let driver = crate::mcp_bulk_auth::McpBulkAuth::new();
let outcomes = driver
.prepare(
servers,
crate::mcp_bulk_auth::BulkAuthMode::Expired,
REAUTH_REDIRECT_URI,
)
.await;
let list: Vec<VmValue> = outcomes
.iter()
.map(|outcome| json_to_vm_value(&crate::mcp_bulk_auth::prepare_outcome_to_json(outcome)))
.collect();
Ok(VmValue::List(std::sync::Arc::new(list)))
}
#[cfg(test)]
mod reauth_tests {
use super::*;
#[test]
fn oauth_servers_from_specs_selects_only_interactive_http() {
let specs = vec![
serde_json::json!({
"name": "Notion",
"transport": "http",
"url": "https://mcp.notion.com/mcp",
}),
serde_json::json!({
"name": "Internal",
"transport": "http",
"url": "https://internal/mcp",
"auth_token": "secret",
}),
serde_json::json!({
"name": "Local",
"transport": "stdio",
"command": "./server",
}),
serde_json::json!({
"name": "ImplicitStdio",
"url": "https://implicit/mcp",
}),
serde_json::json!({ "name": "Bad", "transport": "http", "url": "" }),
];
let servers = oauth_servers_from_specs(&specs);
assert_eq!(servers.len(), 1);
assert_eq!(servers[0].name, "Notion");
assert_eq!(servers[0].server_url, "https://mcp.notion.com/mcp");
}
}