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    tenant: Option<String>,
74    agent: Agent,
75}
76
77impl Client {
78    /// A client for `base_url` (e.g. `http://localhost:5170`).
79    pub fn new(base_url: impl Into<String>) -> Self {
80        let agent: Agent = Agent::config_builder()
81            .http_status_as_error(false)
82            .timeout_global(Some(Duration::from_secs(30 * 60)))
83            .build()
84            .into();
85        Self {
86            base: base_url.into().trim_end_matches('/').to_string(),
87            api_key: None,
88            tenant: None,
89            agent,
90        }
91    }
92
93    /// The credential: a tenant key (`mbk_…`), the platform key (Marbots:ApiKey) or an OIDC access token.
94    pub fn with_api_key(mut self, key: impl Into<String>) -> Self {
95        self.api_key = Some(key.into());
96        self
97    }
98
99    /// Tenant to act in with the platform key or a token (or end the base URL in `/t/<tenant>`).
100    pub fn with_tenant(mut self, tenant: impl Into<String>) -> Self {
101        self.tenant = Some(tenant.into());
102        self
103    }
104
105    /// Adds the credential and tenant headers. OIDC tokens (JWT) go in `Authorization`, keys in `X-Api-Key`.
106    pub(crate) fn auth<B>(&self, mut req: ureq::RequestBuilder<B>) -> ureq::RequestBuilder<B> {
107        if let Some(k) = &self.api_key {
108            req = if k.matches('.').count() == 2 && !k.starts_with("mbk_") {
109                req.header("Authorization", format!("Bearer {k}"))
110            } else {
111                req.header("X-Api-Key", k)
112            };
113        }
114        if let Some(t) = &self.tenant {
115            req = req.header("X-Marbots-Tenant", t);
116        }
117        req
118    }
119
120    pub fn base_url(&self) -> &str {
121        &self.base
122    }
123
124    fn finish(resp: ureq::http::Response<ureq::Body>) -> Result<Vec<u8>> {
125        let status = resp.status().as_u16();
126        let bytes = resp.into_body().read_to_vec()?;
127        if status >= 300 {
128            let text = String::from_utf8_lossy(&bytes).to_string();
129            let message = serde_json::from_str::<Value>(&text)
130                .ok()
131                .and_then(|v| {
132                    v.get("detail")
133                        .or_else(|| v.get("title"))
134                        .and_then(Value::as_str)
135                        .map(str::to_string)
136                })
137                .unwrap_or(text);
138            return Err(Error::Api { status, message });
139        }
140        Ok(bytes)
141    }
142
143    fn decode<T: DeserializeOwned>(bytes: &[u8]) -> Result<T> {
144        Ok(serde_json::from_slice(if bytes.is_empty() {
145            b"null"
146        } else {
147            bytes
148        })?)
149    }
150
151    pub(crate) fn get<T: DeserializeOwned>(&self, path: &str) -> Result<T> {
152        Self::decode(&self.get_raw(path)?)
153    }
154
155    pub(crate) fn get_raw(&self, path: &str) -> Result<Vec<u8>> {
156        let mut req = self.agent.get(format!("{}{}", self.base, path));
157        req = self.auth(req);
158        Self::finish(req.call()?)
159    }
160
161    pub(crate) fn send<B: Serialize, T: DeserializeOwned>(
162        &self,
163        method: &str,
164        path: &str,
165        body: Option<&B>,
166    ) -> Result<T> {
167        let url = format!("{}{}", self.base, path);
168        let resp = match method {
169            "DELETE" => {
170                let mut req = self.agent.delete(url);
171                req = self.auth(req);
172                req.call()?
173            }
174            _ => {
175                let mut req = match method {
176                    "PUT" => self.agent.put(url),
177                    _ => self.agent.post(url),
178                };
179                req = self.auth(req);
180                match body {
181                    Some(b) => req.send_json(b)?,
182                    None => req
183                        .header("Content-Type", "application/json")
184                        .send(&b""[..])?,
185                }
186            }
187        };
188        Self::decode(&Self::finish(resp)?)
189    }
190
191    pub(crate) fn post_bytes(
192        &self,
193        path: &str,
194        body: &[u8],
195        content_type: &str,
196    ) -> Result<Vec<u8>> {
197        let mut req = self
198            .agent
199            .post(format!("{}{}", self.base, path))
200            .header("Content-Type", content_type);
201        req = self.auth(req);
202        Self::finish(req.send(body)?)
203    }
204
205    /// Server information.
206    pub fn system(&self) -> Result<SystemInfo> {
207        self.get("/api/v1/system")
208    }
209
210    /// Machines that run bots.
211    pub fn hosts(&self) -> Result<Vec<HostInfo>> {
212        self.get("/api/v1/hosts")
213    }
214
215    /// Creates a thread with `bot`, sends `text`, waits, and returns the reply text.
216    pub fn chat(&self, bot: &str, text: &str) -> Result<String> {
217        let thread = self.threads().create(bot, None)?;
218        Ok(self
219            .threads()
220            .send(&thread.id, text, true)?
221            .text()
222            .to_string())
223    }
224
225    pub fn bots(&self) -> Bots<'_> {
226        Bots(self)
227    }
228    pub fn templates(&self) -> Templates<'_> {
229        Templates(self)
230    }
231    pub fn models(&self) -> Models<'_> {
232        Models(self)
233    }
234    pub fn threads(&self) -> Threads<'_> {
235        Threads(self)
236    }
237    pub fn tasks(&self) -> Tasks<'_> {
238        Tasks(self)
239    }
240    pub fn approvals(&self) -> Approvals<'_> {
241        Approvals(self)
242    }
243    pub fn skills(&self) -> Skills<'_> {
244        Skills(self)
245    }
246    pub fn mcp(&self) -> Mcp<'_> {
247        Mcp(self)
248    }
249    pub fn schedules(&self) -> Schedules<'_> {
250        Schedules(self)
251    }
252    pub fn memory(&self) -> Memory<'_> {
253        Memory(self)
254    }
255    /// Computers that run bots' tools.
256    pub fn agent_hosts(&self) -> AgentHosts<'_> {
257        AgentHosts(self)
258    }
259    /// Who am I, tenants (platform admins), this tenant's API keys and members (owners).
260    pub fn tenancy(&self) -> Tenancy<'_> {
261        Tenancy(self)
262    }
263    pub fn events(&self) -> Events<'_> {
264        Events(self)
265    }
266}
267
268const NONE: Option<&Value> = None;
269
270/// Bot management.
271pub struct Bots<'a>(&'a Client);
272
273impl Bots<'_> {
274    pub fn list(&self) -> Result<Vec<Bot>> {
275        self.0.get("/api/v1/bots")
276    }
277    pub fn get(&self, id_or_name: &str) -> Result<Bot> {
278        self.0.get(&format!("/api/v1/bots/{}", esc(id_or_name)))
279    }
280    pub fn create(&self, spec: BotSpec) -> Result<Bot> {
281        self.0.send("POST", "/api/v1/bots", Some(&spec))
282    }
283    pub fn update(&self, id: &str, spec: BotSpec) -> Result<Bot> {
284        self.0.send(
285            "PUT",
286            &format!("/api/v1/bots/{}", esc(id)),
287            Some(&spec.with_id(id)),
288        )
289    }
290    pub fn delete(&self, id: &str) -> Result<()> {
291        self.0
292            .send::<Value, Value>("DELETE", &format!("/api/v1/bots/{}", esc(id)), NONE)
293            .map(|_| ())
294    }
295    /// Creates a bot from a gallery template.
296    pub fn hire(&self, template_id: &str, name: Option<&str>) -> Result<Bot> {
297        self.0.send(
298            "POST",
299            &format!("/api/v1/bots/from-template/{}", esc(template_id)),
300            Some(&json!({ "name": name })),
301        )
302    }
303    pub fn get_model(&self, id: &str) -> Result<BotModelInfo> {
304        self.0.get(&format!("/api/v1/bots/{}/model", esc(id)))
305    }
306    /// Sets the bot's model: [`ModelRef::DEFAULT`], [`ModelRef::of`] or a profile name.
307    pub fn set_model(&self, id: &str, model: &str) -> Result<BotModelInfo> {
308        self.0.send(
309            "PUT",
310            &format!("/api/v1/bots/{}/model", esc(id)),
311            Some(&json!({ "model": model })),
312        )
313    }
314    pub fn pause(&self, id: &str) -> Result<()> {
315        self.0
316            .send::<Value, Value>("POST", &format!("/api/v1/bots/{}/pause", esc(id)), NONE)
317            .map(|_| ())
318    }
319    pub fn resume(&self, id: &str) -> Result<()> {
320        self.0
321            .send::<Value, Value>("POST", &format!("/api/v1/bots/{}/resume", esc(id)), NONE)
322            .map(|_| ())
323    }
324    /// Downloads a .marbot package (secrets are never included).
325    pub fn export(&self, id: &str, include_memory: bool) -> Result<Vec<u8>> {
326        self.0.get_raw(&format!(
327            "/api/v1/bots/{}/export?includeMemory={include_memory}",
328            esc(id)
329        ))
330    }
331    pub fn import_package(&self, package: &[u8]) -> Result<Bot> {
332        Client::decode(
333            &self
334                .0
335                .post_bytes("/api/v1/bots/import", package, "application/zip")?,
336        )
337    }
338}
339
340/// Template gallery.
341pub struct Templates<'a>(&'a Client);
342
343impl Templates<'_> {
344    pub fn list(&self, query: &str, category: &str) -> Result<Vec<BotTemplate>> {
345        self.0.get(&format!(
346            "/api/v1/templates?q={}&category={}",
347            esc(query),
348            esc(category)
349        ))
350    }
351    pub fn get(&self, id: &str) -> Result<BotTemplate> {
352        self.0.get(&format!("/api/v1/templates/{}", esc(id)))
353    }
354}
355
356/// Workspace default model and model choices.
357pub struct Models<'a>(&'a Client);
358
359impl Models<'_> {
360    pub fn list(&self) -> Result<ModelCatalog> {
361        self.0.get("/api/v1/models")
362    }
363    /// Changes the model used by every bot whose model is [`ModelRef::DEFAULT`]; returns the new default.
364    pub fn set_default(&self, model: &str) -> Result<String> {
365        let v: Value = self.0.send(
366            "PUT",
367            "/api/v1/models/default",
368            Some(&json!({ "model": model })),
369        )?;
370        Ok(v.get("default")
371            .and_then(Value::as_str)
372            .unwrap_or_default()
373            .to_string())
374    }
375}
376
377/// Conversations.
378pub struct Threads<'a>(&'a Client);
379
380impl Threads<'_> {
381    pub fn list(&self, bot_id: Option<&str>) -> Result<Vec<ChatThread>> {
382        match bot_id {
383            Some(b) => self.0.get(&format!("/api/v1/threads?botId={}", esc(b))),
384            None => self.0.get("/api/v1/threads"),
385        }
386    }
387    pub fn create(&self, bot_id: &str, title: Option<&str>) -> Result<ChatThread> {
388        self.0.send(
389            "POST",
390            "/api/v1/threads",
391            Some(&json!({ "botId": bot_id, "title": title })),
392        )
393    }
394    /// Sends a message; with `wait` the call returns after the bot finished (10 minute limit).
395    pub fn send(&self, thread_id: &str, text: &str, wait: bool) -> Result<SendResult> {
396        self.send_with_timeout(thread_id, text, wait, 600)
397    }
398    pub fn send_with_timeout(
399        &self,
400        thread_id: &str,
401        text: &str,
402        wait: bool,
403        timeout_seconds: u32,
404    ) -> Result<SendResult> {
405        let body = json!({ "text": text, "wait": wait, "timeoutSeconds": timeout_seconds });
406        self.0.send(
407            "POST",
408            &format!("/api/v1/threads/{}/messages", esc(thread_id)),
409            Some(&body),
410        )
411    }
412    pub fn messages(&self, thread_id: &str) -> Result<Vec<ChatMessage>> {
413        self.0
414            .get(&format!("/api/v1/threads/{}/messages", esc(thread_id)))
415    }
416    pub fn files(&self, thread_id: &str) -> Result<Vec<WorkspaceFile>> {
417        self.0
418            .get(&format!("/api/v1/threads/{}/files", esc(thread_id)))
419    }
420    pub fn download(&self, thread_id: &str, path: &str) -> Result<Vec<u8>> {
421        self.0.get_raw(&format!(
422            "/api/v1/threads/{}/files/{}",
423            esc(thread_id),
424            path
425        ))
426    }
427    pub fn delete(&self, thread_id: &str) -> Result<()> {
428        self.0
429            .send::<Value, Value>(
430                "DELETE",
431                &format!("/api/v1/threads/{}", esc(thread_id)),
432                NONE,
433            )
434            .map(|_| ())
435    }
436}
437
438/// Tasks.
439pub struct Tasks<'a>(&'a Client);
440
441impl Tasks<'_> {
442    pub fn list(&self, thread_id: Option<&str>) -> Result<Vec<TaskRecord>> {
443        match thread_id {
444            Some(t) => self.0.get(&format!("/api/v1/tasks?threadId={}", esc(t))),
445            None => self.0.get("/api/v1/tasks"),
446        }
447    }
448    pub fn get(&self, id: &str) -> Result<TaskRecord> {
449        self.0.get(&format!("/api/v1/tasks/{}", esc(id)))
450    }
451    pub fn cancel(&self, id: &str) -> Result<()> {
452        self.0
453            .send::<Value, Value>("POST", &format!("/api/v1/tasks/{}/cancel", esc(id)), NONE)
454            .map(|_| ())
455    }
456}
457
458/// Human-in-the-loop approvals.
459pub struct Approvals<'a>(&'a Client);
460
461impl Approvals<'_> {
462    pub fn pending(&self) -> Result<Vec<ApprovalRequest>> {
463        self.0.get("/api/v1/approvals?state=pending")
464    }
465    pub fn approve(&self, id: &str, scope: ApprovalScope) -> Result<ApprovalRequest> {
466        self.0.send(
467            "POST",
468            &format!("/api/v1/approvals/{}/approve", esc(id)),
469            Some(&json!({ "scope": scope })),
470        )
471    }
472    pub fn reject(&self, id: &str) -> Result<ApprovalRequest> {
473        self.0.send::<Value, _>(
474            "POST",
475            &format!("/api/v1/approvals/{}/reject", esc(id)),
476            NONE,
477        )
478    }
479    /// True when approvals are skipped (dangerous mode).
480    pub fn skip_approvals(&self) -> Result<bool> {
481        let v: Value = self.0.get("/api/v1/system/approvals")?;
482        Ok(v.get("dangerouslySkipApprovals")
483            .and_then(Value::as_bool)
484            .unwrap_or(false))
485    }
486    /// Dangerous, like `--dangerously-skip-permissions`: every action that would ask runs without a human.
487    /// Turning it on also approves everything pending. Actions a bot's profile denies stay denied.
488    pub fn set_skip_approvals(&self, skip: bool) -> Result<bool> {
489        let v: Value = self.0.send(
490            "PUT",
491            "/api/v1/system/approvals",
492            Some(&json!({ "dangerouslySkipApprovals": skip })),
493        )?;
494        Ok(v.get("dangerouslySkipApprovals")
495            .and_then(Value::as_bool)
496            .unwrap_or(false))
497    }
498}
499
500/// SKILL.md packages.
501pub struct Skills<'a>(&'a Client);
502
503impl Skills<'_> {
504    pub fn list(&self) -> Result<Vec<SkillInfo>> {
505        self.0.get("/api/v1/skills")
506    }
507    pub fn install(&self, source: &str) -> Result<Vec<SkillInfo>> {
508        self.0.send(
509            "POST",
510            "/api/v1/skills/install",
511            Some(&json!({ "source": source })),
512        )
513    }
514    /// Learning evaluation: outcomes per skill version and a verdict.
515    pub fn evaluations(&self) -> Result<Vec<SkillEvaluation>> {
516        self.0.get("/api/v1/skills/evaluations")
517    }
518    /// Restores the previous version of an installed skill; returns the restored version.
519    pub fn rollback(&self, name: &str) -> Result<String> {
520        let v: Value = self.0.send(
521            "POST",
522            &format!("/api/v1/skills/{}/rollback", esc(name)),
523            NONE,
524        )?;
525        Ok(v.get("version")
526            .and_then(Value::as_str)
527            .unwrap_or_default()
528            .to_string())
529    }
530    /// Publishes a skill drafted by auto-learn.
531    pub fn promote(&self, name: &str) -> Result<()> {
532        self.0
533            .send::<Value, Value>(
534                "POST",
535                &format!("/api/v1/skills/{}/approve", esc(name)),
536                NONE,
537            )
538            .map(|_| ())
539    }
540    pub fn discard(&self, name: &str) -> Result<()> {
541        self.0
542            .send::<Value, Value>(
543                "POST",
544                &format!("/api/v1/skills/{}/reject", esc(name)),
545                NONE,
546            )
547            .map(|_| ())
548    }
549    pub fn auto_rollback(&self) -> Result<bool> {
550        let v: Value = self.0.get("/api/v1/system/learning")?;
551        Ok(v.get("autoRollbackSkills")
552            .and_then(Value::as_bool)
553            .unwrap_or(false))
554    }
555    pub fn set_auto_rollback(&self, on: bool) -> Result<bool> {
556        let v: Value = self.0.send(
557            "PUT",
558            "/api/v1/system/learning",
559            Some(&json!({ "autoRollbackSkills": on })),
560        )?;
561        Ok(v.get("autoRollbackSkills")
562            .and_then(Value::as_bool)
563            .unwrap_or(false))
564    }
565}
566
567/// Computers that run bots' tools (see docs/en/computers.md).
568pub struct AgentHosts<'a>(&'a Client);
569
570impl AgentHosts<'_> {
571    pub fn list(&self) -> Result<Vec<HostInfo>> {
572        self.0.get("/api/v1/hosts")
573    }
574    /// One-time token for `marbots-host enroll`.
575    pub fn create_enrollment(&self, name: &str, valid_minutes: u32) -> Result<EnrollmentToken> {
576        self.0.send(
577            "POST",
578            "/api/v1/hosts/enrollments",
579            Some(&json!({ "name": name, "validMinutes": valid_minutes })),
580        )
581    }
582    /// Installs marbots-host over SSH; the password/key is used for this call only.
583    pub fn bootstrap(&self, options: &BootstrapOptions) -> Result<BootstrapResult> {
584        self.0
585            .send("POST", "/api/v1/hosts/bootstrap", Some(options))
586    }
587    pub fn disable(&self, id: &str) -> Result<()> {
588        self.0
589            .send::<Value, Value>("POST", &format!("/api/v1/hosts/{}/disable", esc(id)), NONE)
590            .map(|_| ())
591    }
592    pub fn enable(&self, id: &str) -> Result<()> {
593        self.0
594            .send::<Value, Value>("POST", &format!("/api/v1/hosts/{}/enable", esc(id)), NONE)
595            .map(|_| ())
596    }
597    pub fn remove(&self, id: &str) -> Result<()> {
598        self.0
599            .send::<Value, Value>("DELETE", &format!("/api/v1/hosts/{}", esc(id)), NONE)
600            .map(|_| ())
601    }
602}
603
604/// Who am I, tenants and this tenant's keys and members (see docs/en/multi-tenant.md).
605pub struct Tenancy<'a>(&'a Client);
606
607impl Tenancy<'_> {
608    pub fn whoami(&self) -> Result<WhoAmI> {
609        self.0.get("/api/v1/whoami")
610    }
611    pub fn list_tenants(&self) -> Result<Vec<Tenant>> {
612        self.0.get("/api/v1/tenants")
613    }
614    pub fn create_tenant(&self, id: &str, name: Option<&str>) -> Result<Tenant> {
615        self.0.send(
616            "POST",
617            "/api/v1/tenants",
618            Some(&json!({ "id": id, "name": name })),
619        )
620    }
621    pub fn disable_tenant(&self, id: &str) -> Result<Tenant> {
622        self.0.send::<Value, _>(
623            "POST",
624            &format!("/api/v1/tenants/{}/disable", esc(id)),
625            NONE,
626        )
627    }
628    pub fn enable_tenant(&self, id: &str) -> Result<Tenant> {
629        self.0
630            .send::<Value, _>("POST", &format!("/api/v1/tenants/{}/enable", esc(id)), NONE)
631    }
632    /// A key in any tenant (platform admins); the plaintext key is returned once.
633    pub fn create_tenant_key(
634        &self,
635        tenant: &str,
636        name: &str,
637        role: TenantRole,
638    ) -> Result<NewApiKey> {
639        self.0.send(
640            "POST",
641            &format!("/api/v1/tenants/{}/keys", esc(tenant)),
642            Some(&json!({ "name": name, "role": role })),
643        )
644    }
645    pub fn list_keys(&self) -> Result<Vec<ApiKeyInfo>> {
646        self.0.get("/api/v1/tenant/keys")
647    }
648    pub fn create_key(&self, name: &str, role: TenantRole) -> Result<NewApiKey> {
649        self.0.send(
650            "POST",
651            "/api/v1/tenant/keys",
652            Some(&json!({ "name": name, "role": role })),
653        )
654    }
655    pub fn revoke_key(&self, id: &str) -> Result<()> {
656        self.0
657            .send::<Value, Value>("DELETE", &format!("/api/v1/tenant/keys/{}", esc(id)), NONE)
658            .map(|_| ())
659    }
660    pub fn list_members(&self) -> Result<Vec<TenantMember>> {
661        self.0.get("/api/v1/tenant/members")
662    }
663    pub fn set_member(&self, subject: &str, role: TenantRole) -> Result<TenantMember> {
664        self.0.send(
665            "PUT",
666            "/api/v1/tenant/members",
667            Some(&json!({ "subject": subject, "role": role })),
668        )
669    }
670    pub fn remove_member(&self, subject: &str) -> Result<()> {
671        self.0
672            .send::<Value, Value>(
673                "DELETE",
674                &format!("/api/v1/tenant/members/{}", esc(subject)),
675                NONE,
676            )
677            .map(|_| ())
678    }
679}
680
681/// MCP servers.
682pub struct Mcp<'a>(&'a Client);
683
684impl Mcp<'_> {
685    pub fn list(&self) -> Result<Vec<McpServer>> {
686        self.0.get("/api/v1/mcp")
687    }
688    pub fn install(&self, id: &str) -> Result<McpServer> {
689        self.0
690            .send::<Value, _>("POST", &format!("/api/v1/mcp/{}/install", esc(id)), NONE)
691    }
692}
693
694/// Cron and one-off jobs.
695pub struct Schedules<'a>(&'a Client);
696
697impl Schedules<'_> {
698    pub fn list(&self) -> Result<Vec<ScheduleJob>> {
699        self.0.get("/api/v1/schedules")
700    }
701    pub fn create(&self, spec: &ScheduleSpec) -> Result<ScheduleJob> {
702        self.0.send("POST", "/api/v1/schedules", Some(spec))
703    }
704    pub fn delete(&self, id: &str) -> Result<()> {
705        self.0
706            .send::<Value, Value>("DELETE", &format!("/api/v1/schedules/{}", esc(id)), NONE)
707            .map(|_| ())
708    }
709}
710
711/// Long-term memory.
712pub struct Memory<'a>(&'a Client);
713
714impl Memory<'_> {
715    pub fn list(&self, owner: &str) -> Result<Vec<MemoryRecord>> {
716        self.0.get(&format!("/api/v1/memory/{}", esc(owner)))
717    }
718    pub fn remember(&self, owner: &str, content: &str, kind: MemoryKind) -> Result<MemoryRecord> {
719        let body = json!({ "owner": owner, "content": content, "kind": kind, "source": "sdk:rust", "confidence": 1.0 });
720        self.0.send("POST", "/api/v1/memory", Some(&body))
721    }
722}
723
724/// Live events (Server-Sent Events).
725pub struct Events<'a>(&'a Client);
726
727impl Events<'_> {
728    /// Opens the event stream (one thread, or everything when `thread_id` is `None`). Iterate it on a background
729    /// thread if you need to keep working; dropping the iterator closes the connection.
730    pub fn stream(&self, thread_id: Option<&str>) -> Result<EventStream> {
731        let path = match thread_id {
732            Some(t) => format!("/api/v1/threads/{}/events", esc(t)),
733            None => "/api/v1/events".to_string(),
734        };
735        let agent: Agent = Agent::config_builder()
736            .http_status_as_error(false)
737            .build()
738            .into();
739        let mut req = agent
740            .get(format!("{}{}", self.0.base, path))
741            .header("Accept", "text/event-stream");
742        req = self.0.auth(req);
743        let resp = req.call()?;
744        if resp.status().as_u16() >= 300 {
745            return Err(Error::Api {
746                status: resp.status().as_u16(),
747                message: "event stream unavailable".into(),
748            });
749        }
750        Ok(EventStream {
751            lines: Box::new(BufReader::new(resp.into_body().into_reader())),
752        })
753    }
754}
755
756/// Iterator over live events.
757pub struct EventStream {
758    lines: Box<dyn BufRead + Send>,
759}
760
761impl Iterator for EventStream {
762    type Item = Result<AgentEvent>;
763
764    fn next(&mut self) -> Option<Self::Item> {
765        let mut line = String::new();
766        loop {
767            line.clear();
768            match self.lines.read_line(&mut line) {
769                Ok(0) => return None,
770                Ok(_) => {
771                    if let Some(data) = line.trim_end().strip_prefix("data: ") {
772                        return Some(serde_json::from_str(data).map_err(Error::from));
773                    }
774                }
775                Err(e) => return Some(Err(Error::Transport(e.to_string()))),
776            }
777        }
778    }
779}