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    /// Computers that run bots' tools.
241    pub fn agent_hosts(&self) -> AgentHosts<'_> {
242        AgentHosts(self)
243    }
244    pub fn events(&self) -> Events<'_> {
245        Events(self)
246    }
247}
248
249const NONE: Option<&Value> = None;
250
251/// Bot management.
252pub struct Bots<'a>(&'a Client);
253
254impl Bots<'_> {
255    pub fn list(&self) -> Result<Vec<Bot>> {
256        self.0.get("/api/v1/bots")
257    }
258    pub fn get(&self, id_or_name: &str) -> Result<Bot> {
259        self.0.get(&format!("/api/v1/bots/{}", esc(id_or_name)))
260    }
261    pub fn create(&self, spec: BotSpec) -> Result<Bot> {
262        self.0.send("POST", "/api/v1/bots", Some(&spec))
263    }
264    pub fn update(&self, id: &str, spec: BotSpec) -> Result<Bot> {
265        self.0.send(
266            "PUT",
267            &format!("/api/v1/bots/{}", esc(id)),
268            Some(&spec.with_id(id)),
269        )
270    }
271    pub fn delete(&self, id: &str) -> Result<()> {
272        self.0
273            .send::<Value, Value>("DELETE", &format!("/api/v1/bots/{}", esc(id)), NONE)
274            .map(|_| ())
275    }
276    /// Creates a bot from a gallery template.
277    pub fn hire(&self, template_id: &str, name: Option<&str>) -> Result<Bot> {
278        self.0.send(
279            "POST",
280            &format!("/api/v1/bots/from-template/{}", esc(template_id)),
281            Some(&json!({ "name": name })),
282        )
283    }
284    pub fn get_model(&self, id: &str) -> Result<BotModelInfo> {
285        self.0.get(&format!("/api/v1/bots/{}/model", esc(id)))
286    }
287    /// Sets the bot's model: [`ModelRef::DEFAULT`], [`ModelRef::of`] or a profile name.
288    pub fn set_model(&self, id: &str, model: &str) -> Result<BotModelInfo> {
289        self.0.send(
290            "PUT",
291            &format!("/api/v1/bots/{}/model", esc(id)),
292            Some(&json!({ "model": model })),
293        )
294    }
295    pub fn pause(&self, id: &str) -> Result<()> {
296        self.0
297            .send::<Value, Value>("POST", &format!("/api/v1/bots/{}/pause", esc(id)), NONE)
298            .map(|_| ())
299    }
300    pub fn resume(&self, id: &str) -> Result<()> {
301        self.0
302            .send::<Value, Value>("POST", &format!("/api/v1/bots/{}/resume", esc(id)), NONE)
303            .map(|_| ())
304    }
305    /// Downloads a .marbot package (secrets are never included).
306    pub fn export(&self, id: &str, include_memory: bool) -> Result<Vec<u8>> {
307        self.0.get_raw(&format!(
308            "/api/v1/bots/{}/export?includeMemory={include_memory}",
309            esc(id)
310        ))
311    }
312    pub fn import_package(&self, package: &[u8]) -> Result<Bot> {
313        Client::decode(
314            &self
315                .0
316                .post_bytes("/api/v1/bots/import", package, "application/zip")?,
317        )
318    }
319}
320
321/// Template gallery.
322pub struct Templates<'a>(&'a Client);
323
324impl Templates<'_> {
325    pub fn list(&self, query: &str, category: &str) -> Result<Vec<BotTemplate>> {
326        self.0.get(&format!(
327            "/api/v1/templates?q={}&category={}",
328            esc(query),
329            esc(category)
330        ))
331    }
332    pub fn get(&self, id: &str) -> Result<BotTemplate> {
333        self.0.get(&format!("/api/v1/templates/{}", esc(id)))
334    }
335}
336
337/// Workspace default model and model choices.
338pub struct Models<'a>(&'a Client);
339
340impl Models<'_> {
341    pub fn list(&self) -> Result<ModelCatalog> {
342        self.0.get("/api/v1/models")
343    }
344    /// Changes the model used by every bot whose model is [`ModelRef::DEFAULT`]; returns the new default.
345    pub fn set_default(&self, model: &str) -> Result<String> {
346        let v: Value = self.0.send(
347            "PUT",
348            "/api/v1/models/default",
349            Some(&json!({ "model": model })),
350        )?;
351        Ok(v.get("default")
352            .and_then(Value::as_str)
353            .unwrap_or_default()
354            .to_string())
355    }
356}
357
358/// Conversations.
359pub struct Threads<'a>(&'a Client);
360
361impl Threads<'_> {
362    pub fn list(&self, bot_id: Option<&str>) -> Result<Vec<ChatThread>> {
363        match bot_id {
364            Some(b) => self.0.get(&format!("/api/v1/threads?botId={}", esc(b))),
365            None => self.0.get("/api/v1/threads"),
366        }
367    }
368    pub fn create(&self, bot_id: &str, title: Option<&str>) -> Result<ChatThread> {
369        self.0.send(
370            "POST",
371            "/api/v1/threads",
372            Some(&json!({ "botId": bot_id, "title": title })),
373        )
374    }
375    /// Sends a message; with `wait` the call returns after the bot finished (10 minute limit).
376    pub fn send(&self, thread_id: &str, text: &str, wait: bool) -> Result<SendResult> {
377        self.send_with_timeout(thread_id, text, wait, 600)
378    }
379    pub fn send_with_timeout(
380        &self,
381        thread_id: &str,
382        text: &str,
383        wait: bool,
384        timeout_seconds: u32,
385    ) -> Result<SendResult> {
386        let body = json!({ "text": text, "wait": wait, "timeoutSeconds": timeout_seconds });
387        self.0.send(
388            "POST",
389            &format!("/api/v1/threads/{}/messages", esc(thread_id)),
390            Some(&body),
391        )
392    }
393    pub fn messages(&self, thread_id: &str) -> Result<Vec<ChatMessage>> {
394        self.0
395            .get(&format!("/api/v1/threads/{}/messages", esc(thread_id)))
396    }
397    pub fn files(&self, thread_id: &str) -> Result<Vec<WorkspaceFile>> {
398        self.0
399            .get(&format!("/api/v1/threads/{}/files", esc(thread_id)))
400    }
401    pub fn download(&self, thread_id: &str, path: &str) -> Result<Vec<u8>> {
402        self.0.get_raw(&format!(
403            "/api/v1/threads/{}/files/{}",
404            esc(thread_id),
405            path
406        ))
407    }
408    pub fn delete(&self, thread_id: &str) -> Result<()> {
409        self.0
410            .send::<Value, Value>(
411                "DELETE",
412                &format!("/api/v1/threads/{}", esc(thread_id)),
413                NONE,
414            )
415            .map(|_| ())
416    }
417}
418
419/// Tasks.
420pub struct Tasks<'a>(&'a Client);
421
422impl Tasks<'_> {
423    pub fn list(&self, thread_id: Option<&str>) -> Result<Vec<TaskRecord>> {
424        match thread_id {
425            Some(t) => self.0.get(&format!("/api/v1/tasks?threadId={}", esc(t))),
426            None => self.0.get("/api/v1/tasks"),
427        }
428    }
429    pub fn get(&self, id: &str) -> Result<TaskRecord> {
430        self.0.get(&format!("/api/v1/tasks/{}", esc(id)))
431    }
432    pub fn cancel(&self, id: &str) -> Result<()> {
433        self.0
434            .send::<Value, Value>("POST", &format!("/api/v1/tasks/{}/cancel", esc(id)), NONE)
435            .map(|_| ())
436    }
437}
438
439/// Human-in-the-loop approvals.
440pub struct Approvals<'a>(&'a Client);
441
442impl Approvals<'_> {
443    pub fn pending(&self) -> Result<Vec<ApprovalRequest>> {
444        self.0.get("/api/v1/approvals?state=pending")
445    }
446    pub fn approve(&self, id: &str, scope: ApprovalScope) -> Result<ApprovalRequest> {
447        self.0.send(
448            "POST",
449            &format!("/api/v1/approvals/{}/approve", esc(id)),
450            Some(&json!({ "scope": scope })),
451        )
452    }
453    pub fn reject(&self, id: &str) -> Result<ApprovalRequest> {
454        self.0.send::<Value, _>(
455            "POST",
456            &format!("/api/v1/approvals/{}/reject", esc(id)),
457            NONE,
458        )
459    }
460    /// True when approvals are skipped (dangerous mode).
461    pub fn skip_approvals(&self) -> Result<bool> {
462        let v: Value = self.0.get("/api/v1/system/approvals")?;
463        Ok(v.get("dangerouslySkipApprovals")
464            .and_then(Value::as_bool)
465            .unwrap_or(false))
466    }
467    /// Dangerous, like `--dangerously-skip-permissions`: every action that would ask runs without a human.
468    /// Turning it on also approves everything pending. Actions a bot's profile denies stay denied.
469    pub fn set_skip_approvals(&self, skip: bool) -> Result<bool> {
470        let v: Value = self.0.send(
471            "PUT",
472            "/api/v1/system/approvals",
473            Some(&json!({ "dangerouslySkipApprovals": skip })),
474        )?;
475        Ok(v.get("dangerouslySkipApprovals")
476            .and_then(Value::as_bool)
477            .unwrap_or(false))
478    }
479}
480
481/// SKILL.md packages.
482pub struct Skills<'a>(&'a Client);
483
484impl Skills<'_> {
485    pub fn list(&self) -> Result<Vec<SkillInfo>> {
486        self.0.get("/api/v1/skills")
487    }
488    pub fn install(&self, source: &str) -> Result<Vec<SkillInfo>> {
489        self.0.send(
490            "POST",
491            "/api/v1/skills/install",
492            Some(&json!({ "source": source })),
493        )
494    }
495    /// Learning evaluation: outcomes per skill version and a verdict.
496    pub fn evaluations(&self) -> Result<Vec<SkillEvaluation>> {
497        self.0.get("/api/v1/skills/evaluations")
498    }
499    /// Restores the previous version of an installed skill; returns the restored version.
500    pub fn rollback(&self, name: &str) -> Result<String> {
501        let v: Value = self.0.send(
502            "POST",
503            &format!("/api/v1/skills/{}/rollback", esc(name)),
504            NONE,
505        )?;
506        Ok(v.get("version")
507            .and_then(Value::as_str)
508            .unwrap_or_default()
509            .to_string())
510    }
511    /// Publishes a skill drafted by auto-learn.
512    pub fn promote(&self, name: &str) -> Result<()> {
513        self.0
514            .send::<Value, Value>(
515                "POST",
516                &format!("/api/v1/skills/{}/approve", esc(name)),
517                NONE,
518            )
519            .map(|_| ())
520    }
521    pub fn discard(&self, name: &str) -> Result<()> {
522        self.0
523            .send::<Value, Value>(
524                "POST",
525                &format!("/api/v1/skills/{}/reject", esc(name)),
526                NONE,
527            )
528            .map(|_| ())
529    }
530    pub fn auto_rollback(&self) -> Result<bool> {
531        let v: Value = self.0.get("/api/v1/system/learning")?;
532        Ok(v.get("autoRollbackSkills")
533            .and_then(Value::as_bool)
534            .unwrap_or(false))
535    }
536    pub fn set_auto_rollback(&self, on: bool) -> Result<bool> {
537        let v: Value = self.0.send(
538            "PUT",
539            "/api/v1/system/learning",
540            Some(&json!({ "autoRollbackSkills": on })),
541        )?;
542        Ok(v.get("autoRollbackSkills")
543            .and_then(Value::as_bool)
544            .unwrap_or(false))
545    }
546}
547
548/// Computers that run bots' tools (see docs/en/computers.md).
549pub struct AgentHosts<'a>(&'a Client);
550
551impl AgentHosts<'_> {
552    pub fn list(&self) -> Result<Vec<HostInfo>> {
553        self.0.get("/api/v1/hosts")
554    }
555    /// One-time token for `marbots-host enroll`.
556    pub fn create_enrollment(&self, name: &str, valid_minutes: u32) -> Result<EnrollmentToken> {
557        self.0.send(
558            "POST",
559            "/api/v1/hosts/enrollments",
560            Some(&json!({ "name": name, "validMinutes": valid_minutes })),
561        )
562    }
563    /// Installs marbots-host over SSH; the password/key is used for this call only.
564    pub fn bootstrap(&self, options: &BootstrapOptions) -> Result<BootstrapResult> {
565        self.0
566            .send("POST", "/api/v1/hosts/bootstrap", Some(options))
567    }
568    pub fn disable(&self, id: &str) -> Result<()> {
569        self.0
570            .send::<Value, Value>("POST", &format!("/api/v1/hosts/{}/disable", esc(id)), NONE)
571            .map(|_| ())
572    }
573    pub fn enable(&self, id: &str) -> Result<()> {
574        self.0
575            .send::<Value, Value>("POST", &format!("/api/v1/hosts/{}/enable", esc(id)), NONE)
576            .map(|_| ())
577    }
578    pub fn remove(&self, id: &str) -> Result<()> {
579        self.0
580            .send::<Value, Value>("DELETE", &format!("/api/v1/hosts/{}", esc(id)), NONE)
581            .map(|_| ())
582    }
583}
584
585/// MCP servers.
586pub struct Mcp<'a>(&'a Client);
587
588impl Mcp<'_> {
589    pub fn list(&self) -> Result<Vec<McpServer>> {
590        self.0.get("/api/v1/mcp")
591    }
592    pub fn install(&self, id: &str) -> Result<McpServer> {
593        self.0
594            .send::<Value, _>("POST", &format!("/api/v1/mcp/{}/install", esc(id)), NONE)
595    }
596}
597
598/// Cron and one-off jobs.
599pub struct Schedules<'a>(&'a Client);
600
601impl Schedules<'_> {
602    pub fn list(&self) -> Result<Vec<ScheduleJob>> {
603        self.0.get("/api/v1/schedules")
604    }
605    pub fn create(&self, spec: &ScheduleSpec) -> Result<ScheduleJob> {
606        self.0.send("POST", "/api/v1/schedules", Some(spec))
607    }
608    pub fn delete(&self, id: &str) -> Result<()> {
609        self.0
610            .send::<Value, Value>("DELETE", &format!("/api/v1/schedules/{}", esc(id)), NONE)
611            .map(|_| ())
612    }
613}
614
615/// Long-term memory.
616pub struct Memory<'a>(&'a Client);
617
618impl Memory<'_> {
619    pub fn list(&self, owner: &str) -> Result<Vec<MemoryRecord>> {
620        self.0.get(&format!("/api/v1/memory/{}", esc(owner)))
621    }
622    pub fn remember(&self, owner: &str, content: &str, kind: MemoryKind) -> Result<MemoryRecord> {
623        let body = json!({ "owner": owner, "content": content, "kind": kind, "source": "sdk:rust", "confidence": 1.0 });
624        self.0.send("POST", "/api/v1/memory", Some(&body))
625    }
626}
627
628/// Live events (Server-Sent Events).
629pub struct Events<'a>(&'a Client);
630
631impl Events<'_> {
632    /// Opens the event stream (one thread, or everything when `thread_id` is `None`). Iterate it on a background
633    /// thread if you need to keep working; dropping the iterator closes the connection.
634    pub fn stream(&self, thread_id: Option<&str>) -> Result<EventStream> {
635        let path = match thread_id {
636            Some(t) => format!("/api/v1/threads/{}/events", esc(t)),
637            None => "/api/v1/events".to_string(),
638        };
639        let agent: Agent = Agent::config_builder()
640            .http_status_as_error(false)
641            .build()
642            .into();
643        let mut req = agent
644            .get(format!("{}{}", self.0.base, path))
645            .header("Accept", "text/event-stream");
646        if let Some(k) = &self.0.api_key {
647            req = req.header("X-Api-Key", k);
648        }
649        let resp = req.call()?;
650        if resp.status().as_u16() >= 300 {
651            return Err(Error::Api {
652                status: resp.status().as_u16(),
653                message: "event stream unavailable".into(),
654            });
655        }
656        Ok(EventStream {
657            lines: Box::new(BufReader::new(resp.into_body().into_reader())),
658        })
659    }
660}
661
662/// Iterator over live events.
663pub struct EventStream {
664    lines: Box<dyn BufRead + Send>,
665}
666
667impl Iterator for EventStream {
668    type Item = Result<AgentEvent>;
669
670    fn next(&mut self) -> Option<Self::Item> {
671        let mut line = String::new();
672        loop {
673            line.clear();
674            match self.lines.read_line(&mut line) {
675                Ok(0) => return None,
676                Ok(_) => {
677                    if let Some(data) = line.trim_end().strip_prefix("data: ") {
678                        return Some(serde_json::from_str(data).map_err(Error::from));
679                    }
680                }
681                Err(e) => return Some(Err(Error::Transport(e.to_string()))),
682            }
683        }
684    }
685}