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#[derive(Debug)]
13pub enum Error {
14 Api { status: u16, message: String },
16 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 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
53pub 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#[derive(Clone)]
70pub struct Client {
71 base: String,
72 api_key: Option<String>,
73 agent: Agent,
74}
75
76impl Client {
77 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 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 pub fn system(&self) -> Result<SystemInfo> {
192 self.get("/api/v1/system")
193 }
194
195 pub fn hosts(&self) -> Result<Vec<HostInfo>> {
197 self.get("/api/v1/hosts")
198 }
199
200 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 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
251pub 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 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 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 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
321pub 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
337pub struct Models<'a>(&'a Client);
339
340impl Models<'_> {
341 pub fn list(&self) -> Result<ModelCatalog> {
342 self.0.get("/api/v1/models")
343 }
344 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
358pub 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 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
419pub 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
439pub 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 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 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
481pub 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 pub fn evaluations(&self) -> Result<Vec<SkillEvaluation>> {
497 self.0.get("/api/v1/skills/evaluations")
498 }
499 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 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
548pub 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 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 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
585pub 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
598pub 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
615pub 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
628pub struct Events<'a>(&'a Client);
630
631impl Events<'_> {
632 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
662pub 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}