use std::io::{BufRead, BufReader};
use std::time::Duration;
use serde::de::DeserializeOwned;
use serde::Serialize;
use serde_json::{json, Value};
use ureq::Agent;
use crate::types::*;
#[derive(Debug)]
pub enum Error {
Api { status: u16, message: String },
Transport(String),
}
impl std::fmt::Display for Error {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
Error::Api { status, message } => write!(f, "HTTP {status}: {message}"),
Error::Transport(m) => write!(f, "transport error: {m}"),
}
}
}
impl std::error::Error for Error {}
impl Error {
pub fn status(&self) -> Option<u16> {
match self {
Error::Api { status, .. } => Some(*status),
Error::Transport(_) => None,
}
}
}
impl From<ureq::Error> for Error {
fn from(e: ureq::Error) -> Self {
Error::Transport(e.to_string())
}
}
impl From<serde_json::Error> for Error {
fn from(e: serde_json::Error) -> Self {
Error::Transport(e.to_string())
}
}
pub type Result<T> = std::result::Result<T, Error>;
fn esc(s: &str) -> String {
let mut out = String::with_capacity(s.len());
for b in s.bytes() {
if b.is_ascii_alphanumeric() || matches!(b, b'-' | b'_' | b'.' | b'~') {
out.push(b as char);
} else {
out.push_str(&format!("%{b:02X}"));
}
}
out
}
#[derive(Clone)]
pub struct Client {
base: String,
api_key: Option<String>,
tenant: Option<String>,
agent: Agent,
}
impl Client {
pub fn new(base_url: impl Into<String>) -> Self {
let agent: Agent = Agent::config_builder()
.http_status_as_error(false)
.timeout_global(Some(Duration::from_secs(30 * 60)))
.build()
.into();
Self {
base: base_url.into().trim_end_matches('/').to_string(),
api_key: None,
tenant: None,
agent,
}
}
pub fn with_api_key(mut self, key: impl Into<String>) -> Self {
self.api_key = Some(key.into());
self
}
pub fn with_tenant(mut self, tenant: impl Into<String>) -> Self {
self.tenant = Some(tenant.into());
self
}
pub(crate) fn auth<B>(&self, mut req: ureq::RequestBuilder<B>) -> ureq::RequestBuilder<B> {
if let Some(k) = &self.api_key {
req = if k.matches('.').count() == 2 && !k.starts_with("mbk_") {
req.header("Authorization", format!("Bearer {k}"))
} else {
req.header("X-Api-Key", k)
};
}
if let Some(t) = &self.tenant {
req = req.header("X-Marbots-Tenant", t);
}
req
}
pub fn base_url(&self) -> &str {
&self.base
}
fn finish(resp: ureq::http::Response<ureq::Body>) -> Result<Vec<u8>> {
let status = resp.status().as_u16();
let bytes = resp.into_body().read_to_vec()?;
if status >= 300 {
let text = String::from_utf8_lossy(&bytes).to_string();
let message = serde_json::from_str::<Value>(&text)
.ok()
.and_then(|v| {
v.get("detail")
.or_else(|| v.get("title"))
.and_then(Value::as_str)
.map(str::to_string)
})
.unwrap_or(text);
return Err(Error::Api { status, message });
}
Ok(bytes)
}
fn decode<T: DeserializeOwned>(bytes: &[u8]) -> Result<T> {
Ok(serde_json::from_slice(if bytes.is_empty() {
b"null"
} else {
bytes
})?)
}
pub(crate) fn get<T: DeserializeOwned>(&self, path: &str) -> Result<T> {
Self::decode(&self.get_raw(path)?)
}
pub(crate) fn get_raw(&self, path: &str) -> Result<Vec<u8>> {
let mut req = self.agent.get(format!("{}{}", self.base, path));
req = self.auth(req);
Self::finish(req.call()?)
}
pub(crate) fn send<B: Serialize, T: DeserializeOwned>(
&self,
method: &str,
path: &str,
body: Option<&B>,
) -> Result<T> {
let url = format!("{}{}", self.base, path);
let resp = match method {
"DELETE" => {
let mut req = self.agent.delete(url);
req = self.auth(req);
req.call()?
}
_ => {
let mut req = match method {
"PUT" => self.agent.put(url),
_ => self.agent.post(url),
};
req = self.auth(req);
match body {
Some(b) => req.send_json(b)?,
None => req
.header("Content-Type", "application/json")
.send(&b""[..])?,
}
}
};
Self::decode(&Self::finish(resp)?)
}
pub(crate) fn post_bytes(
&self,
path: &str,
body: &[u8],
content_type: &str,
) -> Result<Vec<u8>> {
let mut req = self
.agent
.post(format!("{}{}", self.base, path))
.header("Content-Type", content_type);
req = self.auth(req);
Self::finish(req.send(body)?)
}
pub fn system(&self) -> Result<SystemInfo> {
self.get("/api/v1/system")
}
pub fn hosts(&self) -> Result<Vec<HostInfo>> {
self.get("/api/v1/hosts")
}
pub fn chat(&self, bot: &str, text: &str) -> Result<String> {
let thread = self.threads().create(bot, None)?;
Ok(self
.threads()
.send(&thread.id, text, true)?
.text()
.to_string())
}
pub fn bots(&self) -> Bots<'_> {
Bots(self)
}
pub fn templates(&self) -> Templates<'_> {
Templates(self)
}
pub fn models(&self) -> Models<'_> {
Models(self)
}
pub fn threads(&self) -> Threads<'_> {
Threads(self)
}
pub fn tasks(&self) -> Tasks<'_> {
Tasks(self)
}
pub fn approvals(&self) -> Approvals<'_> {
Approvals(self)
}
pub fn skills(&self) -> Skills<'_> {
Skills(self)
}
pub fn mcp(&self) -> Mcp<'_> {
Mcp(self)
}
pub fn schedules(&self) -> Schedules<'_> {
Schedules(self)
}
pub fn memory(&self) -> Memory<'_> {
Memory(self)
}
pub fn agent_hosts(&self) -> AgentHosts<'_> {
AgentHosts(self)
}
pub fn tenancy(&self) -> Tenancy<'_> {
Tenancy(self)
}
pub fn events(&self) -> Events<'_> {
Events(self)
}
}
const NONE: Option<&Value> = None;
pub struct Bots<'a>(&'a Client);
impl Bots<'_> {
pub fn list(&self) -> Result<Vec<Bot>> {
self.0.get("/api/v1/bots")
}
pub fn get(&self, id_or_name: &str) -> Result<Bot> {
self.0.get(&format!("/api/v1/bots/{}", esc(id_or_name)))
}
pub fn create(&self, spec: BotSpec) -> Result<Bot> {
self.0.send("POST", "/api/v1/bots", Some(&spec))
}
pub fn update(&self, id: &str, spec: BotSpec) -> Result<Bot> {
self.0.send(
"PUT",
&format!("/api/v1/bots/{}", esc(id)),
Some(&spec.with_id(id)),
)
}
pub fn delete(&self, id: &str) -> Result<()> {
self.0
.send::<Value, Value>("DELETE", &format!("/api/v1/bots/{}", esc(id)), NONE)
.map(|_| ())
}
pub fn hire(&self, template_id: &str, name: Option<&str>) -> Result<Bot> {
self.0.send(
"POST",
&format!("/api/v1/bots/from-template/{}", esc(template_id)),
Some(&json!({ "name": name })),
)
}
pub fn get_model(&self, id: &str) -> Result<BotModelInfo> {
self.0.get(&format!("/api/v1/bots/{}/model", esc(id)))
}
pub fn set_model(&self, id: &str, model: &str) -> Result<BotModelInfo> {
self.0.send(
"PUT",
&format!("/api/v1/bots/{}/model", esc(id)),
Some(&json!({ "model": model })),
)
}
pub fn pause(&self, id: &str) -> Result<()> {
self.0
.send::<Value, Value>("POST", &format!("/api/v1/bots/{}/pause", esc(id)), NONE)
.map(|_| ())
}
pub fn resume(&self, id: &str) -> Result<()> {
self.0
.send::<Value, Value>("POST", &format!("/api/v1/bots/{}/resume", esc(id)), NONE)
.map(|_| ())
}
pub fn export(&self, id: &str, include_memory: bool) -> Result<Vec<u8>> {
self.0.get_raw(&format!(
"/api/v1/bots/{}/export?includeMemory={include_memory}",
esc(id)
))
}
pub fn import_package(&self, package: &[u8]) -> Result<Bot> {
Client::decode(
&self
.0
.post_bytes("/api/v1/bots/import", package, "application/zip")?,
)
}
}
pub struct Templates<'a>(&'a Client);
impl Templates<'_> {
pub fn list(&self, query: &str, category: &str) -> Result<Vec<BotTemplate>> {
self.0.get(&format!(
"/api/v1/templates?q={}&category={}",
esc(query),
esc(category)
))
}
pub fn get(&self, id: &str) -> Result<BotTemplate> {
self.0.get(&format!("/api/v1/templates/{}", esc(id)))
}
}
pub struct Models<'a>(&'a Client);
impl Models<'_> {
pub fn list(&self) -> Result<ModelCatalog> {
self.0.get("/api/v1/models")
}
pub fn set_default(&self, model: &str) -> Result<String> {
let v: Value = self.0.send(
"PUT",
"/api/v1/models/default",
Some(&json!({ "model": model })),
)?;
Ok(v.get("default")
.and_then(Value::as_str)
.unwrap_or_default()
.to_string())
}
}
pub struct Threads<'a>(&'a Client);
impl Threads<'_> {
pub fn list(&self, bot_id: Option<&str>) -> Result<Vec<ChatThread>> {
match bot_id {
Some(b) => self.0.get(&format!("/api/v1/threads?botId={}", esc(b))),
None => self.0.get("/api/v1/threads"),
}
}
pub fn create(&self, bot_id: &str, title: Option<&str>) -> Result<ChatThread> {
self.0.send(
"POST",
"/api/v1/threads",
Some(&json!({ "botId": bot_id, "title": title })),
)
}
pub fn send(&self, thread_id: &str, text: &str, wait: bool) -> Result<SendResult> {
self.send_with_timeout(thread_id, text, wait, 600)
}
pub fn send_with_timeout(
&self,
thread_id: &str,
text: &str,
wait: bool,
timeout_seconds: u32,
) -> Result<SendResult> {
let body = json!({ "text": text, "wait": wait, "timeoutSeconds": timeout_seconds });
self.0.send(
"POST",
&format!("/api/v1/threads/{}/messages", esc(thread_id)),
Some(&body),
)
}
pub fn messages(&self, thread_id: &str) -> Result<Vec<ChatMessage>> {
self.0
.get(&format!("/api/v1/threads/{}/messages", esc(thread_id)))
}
pub fn files(&self, thread_id: &str) -> Result<Vec<WorkspaceFile>> {
self.0
.get(&format!("/api/v1/threads/{}/files", esc(thread_id)))
}
pub fn download(&self, thread_id: &str, path: &str) -> Result<Vec<u8>> {
self.0.get_raw(&format!(
"/api/v1/threads/{}/files/{}",
esc(thread_id),
path
))
}
pub fn delete(&self, thread_id: &str) -> Result<()> {
self.0
.send::<Value, Value>(
"DELETE",
&format!("/api/v1/threads/{}", esc(thread_id)),
NONE,
)
.map(|_| ())
}
}
pub struct Tasks<'a>(&'a Client);
impl Tasks<'_> {
pub fn list(&self, thread_id: Option<&str>) -> Result<Vec<TaskRecord>> {
match thread_id {
Some(t) => self.0.get(&format!("/api/v1/tasks?threadId={}", esc(t))),
None => self.0.get("/api/v1/tasks"),
}
}
pub fn get(&self, id: &str) -> Result<TaskRecord> {
self.0.get(&format!("/api/v1/tasks/{}", esc(id)))
}
pub fn cancel(&self, id: &str) -> Result<()> {
self.0
.send::<Value, Value>("POST", &format!("/api/v1/tasks/{}/cancel", esc(id)), NONE)
.map(|_| ())
}
}
pub struct Approvals<'a>(&'a Client);
impl Approvals<'_> {
pub fn pending(&self) -> Result<Vec<ApprovalRequest>> {
self.0.get("/api/v1/approvals?state=pending")
}
pub fn approve(&self, id: &str, scope: ApprovalScope) -> Result<ApprovalRequest> {
self.0.send(
"POST",
&format!("/api/v1/approvals/{}/approve", esc(id)),
Some(&json!({ "scope": scope })),
)
}
pub fn reject(&self, id: &str) -> Result<ApprovalRequest> {
self.0.send::<Value, _>(
"POST",
&format!("/api/v1/approvals/{}/reject", esc(id)),
NONE,
)
}
pub fn skip_approvals(&self) -> Result<bool> {
let v: Value = self.0.get("/api/v1/system/approvals")?;
Ok(v.get("dangerouslySkipApprovals")
.and_then(Value::as_bool)
.unwrap_or(false))
}
pub fn set_skip_approvals(&self, skip: bool) -> Result<bool> {
let v: Value = self.0.send(
"PUT",
"/api/v1/system/approvals",
Some(&json!({ "dangerouslySkipApprovals": skip })),
)?;
Ok(v.get("dangerouslySkipApprovals")
.and_then(Value::as_bool)
.unwrap_or(false))
}
}
pub struct Skills<'a>(&'a Client);
impl Skills<'_> {
pub fn list(&self) -> Result<Vec<SkillInfo>> {
self.0.get("/api/v1/skills")
}
pub fn install(&self, source: &str) -> Result<Vec<SkillInfo>> {
self.0.send(
"POST",
"/api/v1/skills/install",
Some(&json!({ "source": source })),
)
}
pub fn evaluations(&self) -> Result<Vec<SkillEvaluation>> {
self.0.get("/api/v1/skills/evaluations")
}
pub fn rollback(&self, name: &str) -> Result<String> {
let v: Value = self.0.send(
"POST",
&format!("/api/v1/skills/{}/rollback", esc(name)),
NONE,
)?;
Ok(v.get("version")
.and_then(Value::as_str)
.unwrap_or_default()
.to_string())
}
pub fn promote(&self, name: &str) -> Result<()> {
self.0
.send::<Value, Value>(
"POST",
&format!("/api/v1/skills/{}/approve", esc(name)),
NONE,
)
.map(|_| ())
}
pub fn discard(&self, name: &str) -> Result<()> {
self.0
.send::<Value, Value>(
"POST",
&format!("/api/v1/skills/{}/reject", esc(name)),
NONE,
)
.map(|_| ())
}
pub fn auto_rollback(&self) -> Result<bool> {
let v: Value = self.0.get("/api/v1/system/learning")?;
Ok(v.get("autoRollbackSkills")
.and_then(Value::as_bool)
.unwrap_or(false))
}
pub fn set_auto_rollback(&self, on: bool) -> Result<bool> {
let v: Value = self.0.send(
"PUT",
"/api/v1/system/learning",
Some(&json!({ "autoRollbackSkills": on })),
)?;
Ok(v.get("autoRollbackSkills")
.and_then(Value::as_bool)
.unwrap_or(false))
}
}
pub struct AgentHosts<'a>(&'a Client);
impl AgentHosts<'_> {
pub fn list(&self) -> Result<Vec<HostInfo>> {
self.0.get("/api/v1/hosts")
}
pub fn create_enrollment(&self, name: &str, valid_minutes: u32) -> Result<EnrollmentToken> {
self.0.send(
"POST",
"/api/v1/hosts/enrollments",
Some(&json!({ "name": name, "validMinutes": valid_minutes })),
)
}
pub fn bootstrap(&self, options: &BootstrapOptions) -> Result<BootstrapResult> {
self.0
.send("POST", "/api/v1/hosts/bootstrap", Some(options))
}
pub fn disable(&self, id: &str) -> Result<()> {
self.0
.send::<Value, Value>("POST", &format!("/api/v1/hosts/{}/disable", esc(id)), NONE)
.map(|_| ())
}
pub fn enable(&self, id: &str) -> Result<()> {
self.0
.send::<Value, Value>("POST", &format!("/api/v1/hosts/{}/enable", esc(id)), NONE)
.map(|_| ())
}
pub fn remove(&self, id: &str) -> Result<()> {
self.0
.send::<Value, Value>("DELETE", &format!("/api/v1/hosts/{}", esc(id)), NONE)
.map(|_| ())
}
}
pub struct Tenancy<'a>(&'a Client);
impl Tenancy<'_> {
pub fn whoami(&self) -> Result<WhoAmI> {
self.0.get("/api/v1/whoami")
}
pub fn list_tenants(&self) -> Result<Vec<Tenant>> {
self.0.get("/api/v1/tenants")
}
pub fn create_tenant(&self, id: &str, name: Option<&str>) -> Result<Tenant> {
self.0.send(
"POST",
"/api/v1/tenants",
Some(&json!({ "id": id, "name": name })),
)
}
pub fn disable_tenant(&self, id: &str) -> Result<Tenant> {
self.0.send::<Value, _>(
"POST",
&format!("/api/v1/tenants/{}/disable", esc(id)),
NONE,
)
}
pub fn enable_tenant(&self, id: &str) -> Result<Tenant> {
self.0
.send::<Value, _>("POST", &format!("/api/v1/tenants/{}/enable", esc(id)), NONE)
}
pub fn create_tenant_key(
&self,
tenant: &str,
name: &str,
role: TenantRole,
) -> Result<NewApiKey> {
self.0.send(
"POST",
&format!("/api/v1/tenants/{}/keys", esc(tenant)),
Some(&json!({ "name": name, "role": role })),
)
}
pub fn list_keys(&self) -> Result<Vec<ApiKeyInfo>> {
self.0.get("/api/v1/tenant/keys")
}
pub fn create_key(&self, name: &str, role: TenantRole) -> Result<NewApiKey> {
self.0.send(
"POST",
"/api/v1/tenant/keys",
Some(&json!({ "name": name, "role": role })),
)
}
pub fn revoke_key(&self, id: &str) -> Result<()> {
self.0
.send::<Value, Value>("DELETE", &format!("/api/v1/tenant/keys/{}", esc(id)), NONE)
.map(|_| ())
}
pub fn list_members(&self) -> Result<Vec<TenantMember>> {
self.0.get("/api/v1/tenant/members")
}
pub fn set_member(&self, subject: &str, role: TenantRole) -> Result<TenantMember> {
self.0.send(
"PUT",
"/api/v1/tenant/members",
Some(&json!({ "subject": subject, "role": role })),
)
}
pub fn remove_member(&self, subject: &str) -> Result<()> {
self.0
.send::<Value, Value>(
"DELETE",
&format!("/api/v1/tenant/members/{}", esc(subject)),
NONE,
)
.map(|_| ())
}
}
pub struct Mcp<'a>(&'a Client);
impl Mcp<'_> {
pub fn list(&self) -> Result<Vec<McpServer>> {
self.0.get("/api/v1/mcp")
}
pub fn install(&self, id: &str) -> Result<McpServer> {
self.0
.send::<Value, _>("POST", &format!("/api/v1/mcp/{}/install", esc(id)), NONE)
}
}
pub struct Schedules<'a>(&'a Client);
impl Schedules<'_> {
pub fn list(&self) -> Result<Vec<ScheduleJob>> {
self.0.get("/api/v1/schedules")
}
pub fn create(&self, spec: &ScheduleSpec) -> Result<ScheduleJob> {
self.0.send("POST", "/api/v1/schedules", Some(spec))
}
pub fn delete(&self, id: &str) -> Result<()> {
self.0
.send::<Value, Value>("DELETE", &format!("/api/v1/schedules/{}", esc(id)), NONE)
.map(|_| ())
}
}
pub struct Memory<'a>(&'a Client);
impl Memory<'_> {
pub fn list(&self, owner: &str) -> Result<Vec<MemoryRecord>> {
self.0.get(&format!("/api/v1/memory/{}", esc(owner)))
}
pub fn remember(&self, owner: &str, content: &str, kind: MemoryKind) -> Result<MemoryRecord> {
let body = json!({ "owner": owner, "content": content, "kind": kind, "source": "sdk:rust", "confidence": 1.0 });
self.0.send("POST", "/api/v1/memory", Some(&body))
}
}
pub struct Events<'a>(&'a Client);
impl Events<'_> {
pub fn stream(&self, thread_id: Option<&str>) -> Result<EventStream> {
let path = match thread_id {
Some(t) => format!("/api/v1/threads/{}/events", esc(t)),
None => "/api/v1/events".to_string(),
};
let agent: Agent = Agent::config_builder()
.http_status_as_error(false)
.build()
.into();
let mut req = agent
.get(format!("{}{}", self.0.base, path))
.header("Accept", "text/event-stream");
req = self.0.auth(req);
let resp = req.call()?;
if resp.status().as_u16() >= 300 {
return Err(Error::Api {
status: resp.status().as_u16(),
message: "event stream unavailable".into(),
});
}
Ok(EventStream {
lines: Box::new(BufReader::new(resp.into_body().into_reader())),
})
}
}
pub struct EventStream {
lines: Box<dyn BufRead + Send>,
}
impl Iterator for EventStream {
type Item = Result<AgentEvent>;
fn next(&mut self) -> Option<Self::Item> {
let mut line = String::new();
loop {
line.clear();
match self.lines.read_line(&mut line) {
Ok(0) => return None,
Ok(_) => {
if let Some(data) = line.trim_end().strip_prefix("data: ") {
return Some(serde_json::from_str(data).map_err(Error::from));
}
}
Err(e) => return Some(Err(Error::Transport(e.to_string()))),
}
}
}
}