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 events(&self) -> Events<'_> {
241 Events(self)
242 }
243}
244
245const NONE: Option<&Value> = None;
246
247pub 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 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 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 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
317pub 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
333pub struct Models<'a>(&'a Client);
335
336impl Models<'_> {
337 pub fn list(&self) -> Result<ModelCatalog> {
338 self.0.get("/api/v1/models")
339 }
340 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
354pub 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 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
415pub 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
435pub 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 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 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
477pub 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
493pub 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
506pub 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
523pub 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
536pub struct Events<'a>(&'a Client);
538
539impl Events<'_> {
540 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
570pub 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}