Skip to main content

marbots_sdk/
client.rs

1use std::io::{BufRead, BufReader};
2use std::time::Duration;
3
4use serde::de::DeserializeOwned;
5use serde::Serialize;
6use serde_json::{json, Value};
7use ureq::Agent;
8
9use crate::types::*;
10
11/// Errors returned by the SDK.
12#[derive(Debug)]
13pub enum Error {
14    /// The server answered with an error status (message from Problem Details).
15    Api { status: u16, message: String },
16    /// Transport or decoding failure.
17    Transport(String),
18}
19
20impl std::fmt::Display for Error {
21    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
22        match self {
23            Error::Api { status, message } => write!(f, "HTTP {status}: {message}"),
24            Error::Transport(m) => write!(f, "transport error: {m}"),
25        }
26    }
27}
28
29impl std::error::Error for Error {}
30
31impl Error {
32    /// The HTTP status for API errors.
33    pub fn status(&self) -> Option<u16> {
34        match self {
35            Error::Api { status, .. } => Some(*status),
36            Error::Transport(_) => None,
37        }
38    }
39}
40
41impl From<ureq::Error> for Error {
42    fn from(e: ureq::Error) -> Self {
43        Error::Transport(e.to_string())
44    }
45}
46
47impl From<serde_json::Error> for Error {
48    fn from(e: serde_json::Error) -> Self {
49        Error::Transport(e.to_string())
50    }
51}
52
53/// SDK result type.
54pub type Result<T> = std::result::Result<T, Error>;
55
56fn esc(s: &str) -> String {
57    let mut out = String::with_capacity(s.len());
58    for b in s.bytes() {
59        if b.is_ascii_alphanumeric() || matches!(b, b'-' | b'_' | b'.' | b'~') {
60            out.push(b as char);
61        } else {
62            out.push_str(&format!("%{b:02X}"));
63        }
64    }
65    out
66}
67
68/// A synchronous client for a Marbots server. Cheap to clone; safe to share between threads.
69#[derive(Clone)]
70pub struct Client {
71    base: String,
72    api_key: Option<String>,
73    agent: Agent,
74}
75
76impl Client {
77    /// A client for `base_url` (e.g. `http://localhost:5170`).
78    pub fn new(base_url: impl Into<String>) -> Self {
79        let agent: Agent = Agent::config_builder()
80            .http_status_as_error(false)
81            .timeout_global(Some(Duration::from_secs(30 * 60)))
82            .build()
83            .into();
84        Self {
85            base: base_url.into().trim_end_matches('/').to_string(),
86            api_key: None,
87            agent,
88        }
89    }
90
91    /// Sends `X-Api-Key` (needed when the server sets Marbots:ApiKey).
92    pub fn with_api_key(mut self, key: impl Into<String>) -> Self {
93        self.api_key = Some(key.into());
94        self
95    }
96
97    pub fn base_url(&self) -> &str {
98        &self.base
99    }
100
101    fn finish(resp: ureq::http::Response<ureq::Body>) -> Result<Vec<u8>> {
102        let status = resp.status().as_u16();
103        let bytes = resp.into_body().read_to_vec()?;
104        if status >= 300 {
105            let text = String::from_utf8_lossy(&bytes).to_string();
106            let message = serde_json::from_str::<Value>(&text)
107                .ok()
108                .and_then(|v| {
109                    v.get("detail")
110                        .or_else(|| v.get("title"))
111                        .and_then(Value::as_str)
112                        .map(str::to_string)
113                })
114                .unwrap_or(text);
115            return Err(Error::Api { status, message });
116        }
117        Ok(bytes)
118    }
119
120    fn decode<T: DeserializeOwned>(bytes: &[u8]) -> Result<T> {
121        Ok(serde_json::from_slice(if bytes.is_empty() {
122            b"null"
123        } else {
124            bytes
125        })?)
126    }
127
128    pub(crate) fn get<T: DeserializeOwned>(&self, path: &str) -> Result<T> {
129        Self::decode(&self.get_raw(path)?)
130    }
131
132    pub(crate) fn get_raw(&self, path: &str) -> Result<Vec<u8>> {
133        let mut req = self.agent.get(format!("{}{}", self.base, path));
134        if let Some(k) = &self.api_key {
135            req = req.header("X-Api-Key", k);
136        }
137        Self::finish(req.call()?)
138    }
139
140    pub(crate) fn send<B: Serialize, T: DeserializeOwned>(
141        &self,
142        method: &str,
143        path: &str,
144        body: Option<&B>,
145    ) -> Result<T> {
146        let url = format!("{}{}", self.base, path);
147        let resp = match method {
148            "DELETE" => {
149                let mut req = self.agent.delete(url);
150                if let Some(k) = &self.api_key {
151                    req = req.header("X-Api-Key", k);
152                }
153                req.call()?
154            }
155            _ => {
156                let mut req = match method {
157                    "PUT" => self.agent.put(url),
158                    _ => self.agent.post(url),
159                };
160                if let Some(k) = &self.api_key {
161                    req = req.header("X-Api-Key", k);
162                }
163                match body {
164                    Some(b) => req.send_json(b)?,
165                    None => req
166                        .header("Content-Type", "application/json")
167                        .send(&b""[..])?,
168                }
169            }
170        };
171        Self::decode(&Self::finish(resp)?)
172    }
173
174    pub(crate) fn post_bytes(
175        &self,
176        path: &str,
177        body: &[u8],
178        content_type: &str,
179    ) -> Result<Vec<u8>> {
180        let mut req = self
181            .agent
182            .post(format!("{}{}", self.base, path))
183            .header("Content-Type", content_type);
184        if let Some(k) = &self.api_key {
185            req = req.header("X-Api-Key", k);
186        }
187        Self::finish(req.send(body)?)
188    }
189
190    /// Server information.
191    pub fn system(&self) -> Result<SystemInfo> {
192        self.get("/api/v1/system")
193    }
194
195    /// Machines that run bots.
196    pub fn hosts(&self) -> Result<Vec<HostInfo>> {
197        self.get("/api/v1/hosts")
198    }
199
200    /// Creates a thread with `bot`, sends `text`, waits, and returns the reply text.
201    pub fn chat(&self, bot: &str, text: &str) -> Result<String> {
202        let thread = self.threads().create(bot, None)?;
203        Ok(self
204            .threads()
205            .send(&thread.id, text, true)?
206            .text()
207            .to_string())
208    }
209
210    pub fn bots(&self) -> Bots<'_> {
211        Bots(self)
212    }
213    pub fn templates(&self) -> Templates<'_> {
214        Templates(self)
215    }
216    pub fn models(&self) -> Models<'_> {
217        Models(self)
218    }
219    pub fn threads(&self) -> Threads<'_> {
220        Threads(self)
221    }
222    pub fn tasks(&self) -> Tasks<'_> {
223        Tasks(self)
224    }
225    pub fn approvals(&self) -> Approvals<'_> {
226        Approvals(self)
227    }
228    pub fn skills(&self) -> Skills<'_> {
229        Skills(self)
230    }
231    pub fn mcp(&self) -> Mcp<'_> {
232        Mcp(self)
233    }
234    pub fn schedules(&self) -> Schedules<'_> {
235        Schedules(self)
236    }
237    pub fn memory(&self) -> Memory<'_> {
238        Memory(self)
239    }
240    pub fn events(&self) -> Events<'_> {
241        Events(self)
242    }
243}
244
245const NONE: Option<&Value> = None;
246
247/// Bot management.
248pub struct Bots<'a>(&'a Client);
249
250impl Bots<'_> {
251    pub fn list(&self) -> Result<Vec<Bot>> {
252        self.0.get("/api/v1/bots")
253    }
254    pub fn get(&self, id_or_name: &str) -> Result<Bot> {
255        self.0.get(&format!("/api/v1/bots/{}", esc(id_or_name)))
256    }
257    pub fn create(&self, spec: BotSpec) -> Result<Bot> {
258        self.0.send("POST", "/api/v1/bots", Some(&spec))
259    }
260    pub fn update(&self, id: &str, spec: BotSpec) -> Result<Bot> {
261        self.0.send(
262            "PUT",
263            &format!("/api/v1/bots/{}", esc(id)),
264            Some(&spec.with_id(id)),
265        )
266    }
267    pub fn delete(&self, id: &str) -> Result<()> {
268        self.0
269            .send::<Value, Value>("DELETE", &format!("/api/v1/bots/{}", esc(id)), NONE)
270            .map(|_| ())
271    }
272    /// Creates a bot from a gallery template.
273    pub fn hire(&self, template_id: &str, name: Option<&str>) -> Result<Bot> {
274        self.0.send(
275            "POST",
276            &format!("/api/v1/bots/from-template/{}", esc(template_id)),
277            Some(&json!({ "name": name })),
278        )
279    }
280    pub fn get_model(&self, id: &str) -> Result<BotModelInfo> {
281        self.0.get(&format!("/api/v1/bots/{}/model", esc(id)))
282    }
283    /// Sets the bot's model: [`ModelRef::DEFAULT`], [`ModelRef::of`] or a profile name.
284    pub fn set_model(&self, id: &str, model: &str) -> Result<BotModelInfo> {
285        self.0.send(
286            "PUT",
287            &format!("/api/v1/bots/{}/model", esc(id)),
288            Some(&json!({ "model": model })),
289        )
290    }
291    pub fn pause(&self, id: &str) -> Result<()> {
292        self.0
293            .send::<Value, Value>("POST", &format!("/api/v1/bots/{}/pause", esc(id)), NONE)
294            .map(|_| ())
295    }
296    pub fn resume(&self, id: &str) -> Result<()> {
297        self.0
298            .send::<Value, Value>("POST", &format!("/api/v1/bots/{}/resume", esc(id)), NONE)
299            .map(|_| ())
300    }
301    /// Downloads a .marbot package (secrets are never included).
302    pub fn export(&self, id: &str, include_memory: bool) -> Result<Vec<u8>> {
303        self.0.get_raw(&format!(
304            "/api/v1/bots/{}/export?includeMemory={include_memory}",
305            esc(id)
306        ))
307    }
308    pub fn import_package(&self, package: &[u8]) -> Result<Bot> {
309        Client::decode(
310            &self
311                .0
312                .post_bytes("/api/v1/bots/import", package, "application/zip")?,
313        )
314    }
315}
316
317/// Template gallery.
318pub struct Templates<'a>(&'a Client);
319
320impl Templates<'_> {
321    pub fn list(&self, query: &str, category: &str) -> Result<Vec<BotTemplate>> {
322        self.0.get(&format!(
323            "/api/v1/templates?q={}&category={}",
324            esc(query),
325            esc(category)
326        ))
327    }
328    pub fn get(&self, id: &str) -> Result<BotTemplate> {
329        self.0.get(&format!("/api/v1/templates/{}", esc(id)))
330    }
331}
332
333/// Workspace default model and model choices.
334pub struct Models<'a>(&'a Client);
335
336impl Models<'_> {
337    pub fn list(&self) -> Result<ModelCatalog> {
338        self.0.get("/api/v1/models")
339    }
340    /// Changes the model used by every bot whose model is [`ModelRef::DEFAULT`]; returns the new default.
341    pub fn set_default(&self, model: &str) -> Result<String> {
342        let v: Value = self.0.send(
343            "PUT",
344            "/api/v1/models/default",
345            Some(&json!({ "model": model })),
346        )?;
347        Ok(v.get("default")
348            .and_then(Value::as_str)
349            .unwrap_or_default()
350            .to_string())
351    }
352}
353
354/// Conversations.
355pub struct Threads<'a>(&'a Client);
356
357impl Threads<'_> {
358    pub fn list(&self, bot_id: Option<&str>) -> Result<Vec<ChatThread>> {
359        match bot_id {
360            Some(b) => self.0.get(&format!("/api/v1/threads?botId={}", esc(b))),
361            None => self.0.get("/api/v1/threads"),
362        }
363    }
364    pub fn create(&self, bot_id: &str, title: Option<&str>) -> Result<ChatThread> {
365        self.0.send(
366            "POST",
367            "/api/v1/threads",
368            Some(&json!({ "botId": bot_id, "title": title })),
369        )
370    }
371    /// Sends a message; with `wait` the call returns after the bot finished (10 minute limit).
372    pub fn send(&self, thread_id: &str, text: &str, wait: bool) -> Result<SendResult> {
373        self.send_with_timeout(thread_id, text, wait, 600)
374    }
375    pub fn send_with_timeout(
376        &self,
377        thread_id: &str,
378        text: &str,
379        wait: bool,
380        timeout_seconds: u32,
381    ) -> Result<SendResult> {
382        let body = json!({ "text": text, "wait": wait, "timeoutSeconds": timeout_seconds });
383        self.0.send(
384            "POST",
385            &format!("/api/v1/threads/{}/messages", esc(thread_id)),
386            Some(&body),
387        )
388    }
389    pub fn messages(&self, thread_id: &str) -> Result<Vec<ChatMessage>> {
390        self.0
391            .get(&format!("/api/v1/threads/{}/messages", esc(thread_id)))
392    }
393    pub fn files(&self, thread_id: &str) -> Result<Vec<WorkspaceFile>> {
394        self.0
395            .get(&format!("/api/v1/threads/{}/files", esc(thread_id)))
396    }
397    pub fn download(&self, thread_id: &str, path: &str) -> Result<Vec<u8>> {
398        self.0.get_raw(&format!(
399            "/api/v1/threads/{}/files/{}",
400            esc(thread_id),
401            path
402        ))
403    }
404    pub fn delete(&self, thread_id: &str) -> Result<()> {
405        self.0
406            .send::<Value, Value>(
407                "DELETE",
408                &format!("/api/v1/threads/{}", esc(thread_id)),
409                NONE,
410            )
411            .map(|_| ())
412    }
413}
414
415/// Tasks.
416pub struct Tasks<'a>(&'a Client);
417
418impl Tasks<'_> {
419    pub fn list(&self, thread_id: Option<&str>) -> Result<Vec<TaskRecord>> {
420        match thread_id {
421            Some(t) => self.0.get(&format!("/api/v1/tasks?threadId={}", esc(t))),
422            None => self.0.get("/api/v1/tasks"),
423        }
424    }
425    pub fn get(&self, id: &str) -> Result<TaskRecord> {
426        self.0.get(&format!("/api/v1/tasks/{}", esc(id)))
427    }
428    pub fn cancel(&self, id: &str) -> Result<()> {
429        self.0
430            .send::<Value, Value>("POST", &format!("/api/v1/tasks/{}/cancel", esc(id)), NONE)
431            .map(|_| ())
432    }
433}
434
435/// Human-in-the-loop approvals.
436pub struct Approvals<'a>(&'a Client);
437
438impl Approvals<'_> {
439    pub fn pending(&self) -> Result<Vec<ApprovalRequest>> {
440        self.0.get("/api/v1/approvals?state=pending")
441    }
442    pub fn approve(&self, id: &str, scope: ApprovalScope) -> Result<ApprovalRequest> {
443        self.0.send(
444            "POST",
445            &format!("/api/v1/approvals/{}/approve", esc(id)),
446            Some(&json!({ "scope": scope })),
447        )
448    }
449    pub fn reject(&self, id: &str) -> Result<ApprovalRequest> {
450        self.0.send::<Value, _>(
451            "POST",
452            &format!("/api/v1/approvals/{}/reject", esc(id)),
453            NONE,
454        )
455    }
456    /// True when approvals are skipped (dangerous mode).
457    pub fn skip_approvals(&self) -> Result<bool> {
458        let v: Value = self.0.get("/api/v1/system/approvals")?;
459        Ok(v.get("dangerouslySkipApprovals")
460            .and_then(Value::as_bool)
461            .unwrap_or(false))
462    }
463    /// Dangerous, like `--dangerously-skip-permissions`: every action that would ask runs without a human.
464    /// Turning it on also approves everything pending. Actions a bot's profile denies stay denied.
465    pub fn set_skip_approvals(&self, skip: bool) -> Result<bool> {
466        let v: Value = self.0.send(
467            "PUT",
468            "/api/v1/system/approvals",
469            Some(&json!({ "dangerouslySkipApprovals": skip })),
470        )?;
471        Ok(v.get("dangerouslySkipApprovals")
472            .and_then(Value::as_bool)
473            .unwrap_or(false))
474    }
475}
476
477/// SKILL.md packages.
478pub struct Skills<'a>(&'a Client);
479
480impl Skills<'_> {
481    pub fn list(&self) -> Result<Vec<SkillInfo>> {
482        self.0.get("/api/v1/skills")
483    }
484    pub fn install(&self, source: &str) -> Result<Vec<SkillInfo>> {
485        self.0.send(
486            "POST",
487            "/api/v1/skills/install",
488            Some(&json!({ "source": source })),
489        )
490    }
491}
492
493/// MCP servers.
494pub struct Mcp<'a>(&'a Client);
495
496impl Mcp<'_> {
497    pub fn list(&self) -> Result<Vec<McpServer>> {
498        self.0.get("/api/v1/mcp")
499    }
500    pub fn install(&self, id: &str) -> Result<McpServer> {
501        self.0
502            .send::<Value, _>("POST", &format!("/api/v1/mcp/{}/install", esc(id)), NONE)
503    }
504}
505
506/// Cron and one-off jobs.
507pub struct Schedules<'a>(&'a Client);
508
509impl Schedules<'_> {
510    pub fn list(&self) -> Result<Vec<ScheduleJob>> {
511        self.0.get("/api/v1/schedules")
512    }
513    pub fn create(&self, spec: &ScheduleSpec) -> Result<ScheduleJob> {
514        self.0.send("POST", "/api/v1/schedules", Some(spec))
515    }
516    pub fn delete(&self, id: &str) -> Result<()> {
517        self.0
518            .send::<Value, Value>("DELETE", &format!("/api/v1/schedules/{}", esc(id)), NONE)
519            .map(|_| ())
520    }
521}
522
523/// Long-term memory.
524pub struct Memory<'a>(&'a Client);
525
526impl Memory<'_> {
527    pub fn list(&self, owner: &str) -> Result<Vec<MemoryRecord>> {
528        self.0.get(&format!("/api/v1/memory/{}", esc(owner)))
529    }
530    pub fn remember(&self, owner: &str, content: &str, kind: MemoryKind) -> Result<MemoryRecord> {
531        let body = json!({ "owner": owner, "content": content, "kind": kind, "source": "sdk:rust", "confidence": 1.0 });
532        self.0.send("POST", "/api/v1/memory", Some(&body))
533    }
534}
535
536/// Live events (Server-Sent Events).
537pub struct Events<'a>(&'a Client);
538
539impl Events<'_> {
540    /// Opens the event stream (one thread, or everything when `thread_id` is `None`). Iterate it on a background
541    /// thread if you need to keep working; dropping the iterator closes the connection.
542    pub fn stream(&self, thread_id: Option<&str>) -> Result<EventStream> {
543        let path = match thread_id {
544            Some(t) => format!("/api/v1/threads/{}/events", esc(t)),
545            None => "/api/v1/events".to_string(),
546        };
547        let agent: Agent = Agent::config_builder()
548            .http_status_as_error(false)
549            .build()
550            .into();
551        let mut req = agent
552            .get(format!("{}{}", self.0.base, path))
553            .header("Accept", "text/event-stream");
554        if let Some(k) = &self.0.api_key {
555            req = req.header("X-Api-Key", k);
556        }
557        let resp = req.call()?;
558        if resp.status().as_u16() >= 300 {
559            return Err(Error::Api {
560                status: resp.status().as_u16(),
561                message: "event stream unavailable".into(),
562            });
563        }
564        Ok(EventStream {
565            lines: Box::new(BufReader::new(resp.into_body().into_reader())),
566        })
567    }
568}
569
570/// Iterator over live events.
571pub struct EventStream {
572    lines: Box<dyn BufRead + Send>,
573}
574
575impl Iterator for EventStream {
576    type Item = Result<AgentEvent>;
577
578    fn next(&mut self) -> Option<Self::Item> {
579        let mut line = String::new();
580        loop {
581            line.clear();
582            match self.lines.read_line(&mut line) {
583                Ok(0) => return None,
584                Ok(_) => {
585                    if let Some(data) = line.trim_end().strip_prefix("data: ") {
586                        return Some(serde_json::from_str(data).map_err(Error::from));
587                    }
588                }
589                Err(e) => return Some(Err(Error::Transport(e.to_string()))),
590            }
591        }
592    }
593}