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 tenant: Option<String>,
74 agent: Agent,
75}
76
77impl Client {
78 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 pub fn with_api_key(mut self, key: impl Into<String>) -> Self {
95 self.api_key = Some(key.into());
96 self
97 }
98
99 pub fn with_tenant(mut self, tenant: impl Into<String>) -> Self {
101 self.tenant = Some(tenant.into());
102 self
103 }
104
105 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 pub fn system(&self) -> Result<SystemInfo> {
207 self.get("/api/v1/system")
208 }
209
210 pub fn hosts(&self) -> Result<Vec<HostInfo>> {
212 self.get("/api/v1/hosts")
213 }
214
215 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 pub fn agent_hosts(&self) -> AgentHosts<'_> {
257 AgentHosts(self)
258 }
259 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
270pub 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 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 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 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
340pub 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
356pub struct Models<'a>(&'a Client);
358
359impl Models<'_> {
360 pub fn list(&self) -> Result<ModelCatalog> {
361 self.0.get("/api/v1/models")
362 }
363 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
377pub 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 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
438pub 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
458pub 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 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 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
500pub 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 pub fn evaluations(&self) -> Result<Vec<SkillEvaluation>> {
516 self.0.get("/api/v1/skills/evaluations")
517 }
518 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 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
567pub 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 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 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
604pub 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 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
681pub 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
694pub 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
711pub 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
724pub struct Events<'a>(&'a Client);
726
727impl Events<'_> {
728 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
756pub 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}