use anyhow::{Context, bail};
use clap::{Args, Subcommand};
use rmcp::{
ServiceExt,
model::{CallToolRequestParams, ClientInfo},
transport::{
StreamableHttpClientTransport, streamable_http_client::StreamableHttpClientTransportConfig,
},
};
use serde_json::{Value, json};
#[derive(Args)]
pub struct ClientArgs {
#[arg(long, env = "BUS_URL", default_value = "http://localhost:8787/mcp")]
pub url: String,
#[arg(long, env = "BUS_TOKEN", hide_env_values = true)]
pub token: String,
#[arg(long, global = true)]
pub json: bool,
#[command(subcommand)]
pub command: ClientCmd,
}
#[derive(Subcommand)]
pub enum ClientCmd {
Whoami,
Tools,
Send {
#[arg(long, conflicts_with = "to")]
channel: Option<String>,
#[arg(long)]
to: Option<String>,
#[arg(long)]
body: String,
#[arg(long)]
reply_to: Option<i64>,
#[arg(long)]
file: Vec<std::path::PathBuf>,
},
Attach {
task: String,
#[arg(long)]
file: std::path::PathBuf,
#[arg(long)]
content_type: Option<String>,
},
Download {
id: i64,
#[arg(long)]
out: Option<std::path::PathBuf>,
},
Ask {
to: String,
question: Option<String>,
#[arg(long)]
timeout_seconds: Option<i64>,
#[arg(long)]
resume_id: Option<i64>,
},
Read {
#[arg(long, default_value = "all")]
scope: String,
#[arg(long)]
history: bool,
#[arg(long, default_value_t = 50)]
limit: i64,
},
Search {
query: String,
#[arg(long, default_value_t = 50)]
limit: i64,
},
Channels,
ChannelCreate {
name: String,
#[arg(long)]
topic: Option<String>,
},
Agents {
#[arg(long)]
online: bool,
},
Beat {
#[arg(long)]
status: Option<String>,
#[arg(long)]
repo: Option<String>,
#[arg(long)]
branch: Option<String>,
#[arg(long)]
activity: Option<String>,
#[arg(long)]
ttl_seconds: Option<i64>,
},
Tasks {
#[arg(long)]
status: Option<String>,
#[arg(long)]
mine: bool,
},
#[command(subcommand)]
Task(TaskCmd),
Notes {
#[arg(long)]
scope: Option<String>,
#[arg(long)]
tag: Option<String>,
},
#[command(subcommand)]
Note(NoteCmd),
Wait {
#[arg(long)]
timeout_seconds: Option<i64>,
#[arg(long, value_delimiter = ',')]
kinds: Vec<String>,
},
#[command(subcommand)]
Lock(LockCmd),
Digest {
#[arg(long, default_value_t = 24)]
hours: i64,
},
Call {
tool: String,
#[arg(long, default_value = "{}")]
args: String,
},
}
#[derive(Subcommand)]
pub enum LockCmd {
Acquire {
name: String,
#[arg(long)]
ttl_seconds: Option<i64>,
#[arg(long)]
purpose: Option<String>,
},
Release {
name: String,
},
List,
}
#[derive(Subcommand)]
pub enum TaskCmd {
Create {
key: String,
#[arg(long)]
title: String,
#[arg(long)]
description: Option<String>,
#[arg(long, value_delimiter = ',')]
depends_on: Vec<String>,
},
Show {
key: String,
},
Claim {
key: String,
#[arg(long)]
lease_seconds: Option<i64>,
},
Next {
#[arg(long)]
lease_seconds: Option<i64>,
},
Renew {
key: String,
#[arg(long)]
lease_seconds: Option<i64>,
},
Release {
key: String,
},
Done {
key: String,
#[arg(long)]
result: Option<String>,
},
}
#[derive(Subcommand)]
pub enum NoteCmd {
Get {
key: String,
#[arg(long)]
scope: Option<String>,
},
Set {
key: String,
#[arg(long)]
value: String,
#[arg(long)]
scope: Option<String>,
#[arg(long, value_delimiter = ',')]
tags: Vec<String>,
},
Rm {
key: String,
#[arg(long)]
scope: Option<String>,
},
Search {
query: String,
#[arg(long)]
scope: Option<String>,
},
}
pub mod mapping;
pub mod render;
use mapping::to_call;
use render::render;
pub async fn run(args: ClientArgs) -> anyhow::Result<()> {
let mut config = StreamableHttpClientTransportConfig::with_uri(args.url.clone());
config.auth_header = Some(args.token.clone());
config.allow_stateless = true;
let transport = StreamableHttpClientTransport::from_config(config);
let client = ClientInfo::default()
.serve(transport)
.await
.context("could not connect to the bus (check --url and --token)")?;
let outcome = run_command(&client, &args).await;
let _ = client.cancel().await;
outcome
}
async fn run_command(
client: &rmcp::service::RunningService<rmcp::RoleClient, ClientInfo>,
args: &ClientArgs,
) -> anyhow::Result<()> {
if matches!(args.command, ClientCmd::Tools) {
let tools = client.list_all_tools().await?;
if args.json {
println!("{}", serde_json::to_string_pretty(&tools)?);
} else {
for tool in tools {
println!(
"{:<20} {}",
tool.name,
tool.description.as_deref().unwrap_or_default().trim()
);
}
}
return Ok(());
}
let (tool, call_args) = match &args.command {
ClientCmd::Call { tool, args: raw } => {
let parsed: Value = serde_json::from_str(raw)
.with_context(|| format!("--args is not valid JSON: {raw}"))?;
if !parsed.is_object() {
bail!("--args must be a JSON object");
}
(tool.clone(), parsed)
}
other => {
let Some((tool, call_args)) = to_call(other)? else {
bail!("this subcommand has no MCP tool mapping yet");
};
(tool.to_string(), call_args)
}
};
let arguments: serde_json::Map<String, Value> =
serde_json::from_value(call_args).context("arguments did not form a JSON object")?;
let result = client
.call_tool(CallToolRequestParams::new(tool.clone()).with_arguments(arguments))
.await
.map_err(|e| anyhow::anyhow!("{tool} failed: {e}"))?;
if result.is_error == Some(true) {
bail!("{tool} returned an error: {:?}", result.content);
}
let value = result
.structured_content
.clone()
.unwrap_or_else(|| json!({ "ok": true }));
if args.json {
println!("{}", serde_json::to_string_pretty(&value)?);
} else {
render(&args.command, &value)?;
}
Ok(())
}