1use serde_json::{json, Map, Value};
9
10use crate::client::Client;
11use crate::errors::Result;
12use crate::protocol::{new_disconnect, new_request_with_id, new_subscribe};
13
14fn opt_request_id(request_id: &[&str]) -> String {
17 request_id
18 .iter()
19 .copied()
20 .find(|s| !s.is_empty())
21 .unwrap_or_default()
22 .to_string()
23}
24
25impl Client {
26 pub async fn send_command(&self, cmd: &str) -> Result<()> {
28 let mut params = Map::new();
29 params.insert("cmd".into(), json!(cmd));
30 self.notify("slash_command", params).await
31 }
32
33 pub async fn send_detach(&self) -> Result<()> {
35 self.send_envelope(new_disconnect()).await
36 }
37
38 pub async fn send_daemon_status(&self, request_id: &[&str]) -> Result<()> {
40 let rid = opt_request_id(request_id);
41 self.send_envelope(new_request_with_id("daemon_status", Map::new(), rid))
42 .await
43 }
44
45 pub async fn send_daemon_shutdown(&self, request_id: &[&str]) -> Result<()> {
47 let rid = opt_request_id(request_id);
48 self.send_envelope(new_request_with_id("daemon_shutdown", Map::new(), rid))
49 .await
50 }
51
52 pub async fn send_config_get(&self, section: &str, request_id: &[&str]) -> Result<()> {
54 let rid = opt_request_id(request_id);
55 let mut params = Map::new();
56 params.insert("section".into(), json!(section));
57 self.send_envelope(new_request_with_id("config_get", params, rid))
58 .await
59 }
60
61 pub async fn send_config_reload(&self, request_id: &[&str]) -> Result<()> {
63 let rid = opt_request_id(request_id);
64 self.send_envelope(new_request_with_id("config_reload", Map::new(), rid))
65 .await
66 }
67
68 pub async fn send_skills_list(&self, request_id: &[&str]) -> Result<()> {
70 let rid = opt_request_id(request_id);
71 self.send_envelope(new_request_with_id("skills_list", Map::new(), rid))
72 .await
73 }
74
75 pub async fn send_models_list(&self, request_id: &[&str]) -> Result<()> {
77 let rid = opt_request_id(request_id);
78 self.send_envelope(new_request_with_id("models_list", Map::new(), rid))
79 .await
80 }
81
82 pub async fn send_invoke_skill(
84 &self,
85 skill: &str,
86 args: &str,
87 request_id: &[&str],
88 ) -> Result<()> {
89 let rid = opt_request_id(request_id);
90 let mut params = Map::new();
91 params.insert("skill".into(), json!(skill));
92 if !args.is_empty() {
93 params.insert("args".into(), json!(args));
94 }
95 self.send_envelope(new_request_with_id("invoke_skill", params, rid))
96 .await
97 }
98
99 pub async fn send_mcp_status(&self, request_id: &[&str]) -> Result<()> {
101 let rid = opt_request_id(request_id);
102 self.send_envelope(new_request_with_id("mcp_status", Map::new(), rid))
103 .await
104 }
105
106 pub async fn send_loop_list(
108 &self,
109 filter: Option<Map<String, Value>>,
110 limit: u32,
111 request_id: &[&str],
112 ) -> Result<()> {
113 let rid = opt_request_id(request_id);
114 let mut params = Map::new();
115 if let Some(f) = filter {
116 params.insert("filter".into(), Value::Object(f));
117 }
118 if limit > 0 {
119 params.insert("limit".into(), json!(limit));
120 }
121 self.send_envelope(new_request_with_id("loop_list", params, rid))
122 .await
123 }
124
125 pub async fn send_loop_get(
127 &self,
128 loop_id: &str,
129 verbose: bool,
130 request_id: &[&str],
131 ) -> Result<()> {
132 let rid = opt_request_id(request_id);
133 let mut params = Map::new();
134 params.insert("loop_id".into(), json!(loop_id));
135 if verbose {
136 params.insert("verbose".into(), json!(true));
137 }
138 self.send_envelope(new_request_with_id("loop_get", params, rid))
139 .await
140 }
141
142 pub async fn send_loop_tree(
144 &self,
145 loop_id: &str,
146 format: &str,
147 request_id: &[&str],
148 ) -> Result<()> {
149 let rid = opt_request_id(request_id);
150 let mut params = Map::new();
151 params.insert("loop_id".into(), json!(loop_id));
152 if !format.is_empty() {
153 params.insert("format".into(), json!(format));
154 }
155 self.send_envelope(new_request_with_id("loop_tree", params, rid))
156 .await
157 }
158
159 pub async fn send_loop_prune(
161 &self,
162 loop_id: &str,
163 keep_latest: u32,
164 request_id: &[&str],
165 ) -> Result<()> {
166 let rid = opt_request_id(request_id);
167 let mut params = Map::new();
168 params.insert("loop_id".into(), json!(loop_id));
169 if keep_latest > 0 {
170 params.insert("keep_latest".into(), json!(keep_latest));
171 }
172 self.send_envelope(new_request_with_id("loop_prune", params, rid))
173 .await
174 }
175
176 pub async fn send_loop_delete(&self, loop_id: &str, request_id: &[&str]) -> Result<()> {
178 let rid = opt_request_id(request_id);
179 let mut params = Map::new();
180 params.insert("loop_id".into(), json!(loop_id));
181 self.send_envelope(new_request_with_id("loop_delete", params, rid))
182 .await
183 }
184
185 pub async fn send_loop_reattach(&self, loop_id: &str, request_id: &[&str]) -> Result<()> {
187 let rid = opt_request_id(request_id);
188 let mut params = Map::new();
189 params.insert("loop_id".into(), json!(loop_id));
190 self.send_envelope(new_request_with_id("loop_reattach", params, rid))
191 .await
192 }
193
194 pub async fn send_loop_subscribe(
198 &self,
199 loop_id: &str,
200 wire_tier: &str,
201 stream_delivery: &str,
202 request_id: &[&str],
203 ) -> Result<()> {
204 let rid = opt_request_id(request_id);
205 let mut params = Map::new();
206 params.insert("loop_id".into(), json!(loop_id));
207 if !wire_tier.is_empty() {
208 params.insert("wire_tier".into(), json!(wire_tier));
209 }
210 if !stream_delivery.is_empty() {
211 params.insert("stream_delivery".into(), json!(stream_delivery));
212 }
213 let env = if rid.is_empty() {
214 new_subscribe("loop_events", params)
215 } else {
216 let mut e = new_subscribe("loop_events", params);
218 e.id = Some(rid);
219 e
220 };
221 self.send_envelope(env).await
222 }
223
224 pub async fn send_loop_detach(&self, loop_id: &str, request_id: &[&str]) -> Result<()> {
228 let rid = opt_request_id(request_id);
229 let id = if rid.is_empty() {
230 loop_id.to_string()
231 } else {
232 rid
233 };
234 self.unsubscribe(&id).await
235 }
236
237 pub async fn send_loop_new(
239 &self,
240 client_workspace: &str,
241 user_id: &str,
242 client_workspace_id: &str,
243 is_ephemeral: bool,
244 request_id: &[&str],
245 ) -> Result<()> {
246 let rid = opt_request_id(request_id);
247 let mut params = Map::new();
248 if !client_workspace.is_empty() {
249 params.insert("client_workspace".into(), json!(client_workspace));
250 }
251 if !user_id.is_empty() {
252 params.insert("user_id".into(), json!(user_id));
253 }
254 if !client_workspace_id.is_empty() {
255 params.insert("client_workspace_id".into(), json!(client_workspace_id));
256 }
257 if is_ephemeral {
258 params.insert("is_ephemeral".into(), json!(true));
259 }
260 self.send_envelope(new_request_with_id("loop_new", params, rid))
261 .await
262 }
263
264 pub async fn send_loop_input(&self, loop_id: &str, content: &str) -> Result<()> {
266 let mut params = Map::new();
267 params.insert("loop_id".into(), json!(loop_id));
268 params.insert("content".into(), json!(content));
269 self.notify("loop_input", params).await
270 }
271
272 pub async fn send_loop_messages(
274 &self,
275 loop_id: &str,
276 limit: u32,
277 offset: u32,
278 include_events: bool,
279 request_id: &[&str],
280 ) -> Result<()> {
281 let rid = opt_request_id(request_id);
282 let mut params = Map::new();
283 params.insert("loop_id".into(), json!(loop_id));
284 if limit > 0 {
285 params.insert("limit".into(), json!(limit));
286 }
287 if offset > 0 {
288 params.insert("offset".into(), json!(offset));
289 }
290 if include_events {
291 params.insert("include_events".into(), json!(true));
292 }
293 self.send_envelope(new_request_with_id("loop_messages", params, rid))
294 .await
295 }
296
297 pub async fn send_loop_state_get(&self, loop_id: &str, request_id: &[&str]) -> Result<()> {
299 let rid = opt_request_id(request_id);
300 let mut params = Map::new();
301 params.insert("loop_id".into(), json!(loop_id));
302 self.send_envelope(new_request_with_id("loop_state_get", params, rid))
303 .await
304 }
305
306 pub async fn send_loop_state_update(
308 &self,
309 loop_id: &str,
310 values: Map<String, Value>,
311 as_node: &str,
312 request_id: &[&str],
313 ) -> Result<()> {
314 let rid = opt_request_id(request_id);
315 let mut params = Map::new();
316 params.insert("loop_id".into(), json!(loop_id));
317 params.insert("values".into(), Value::Object(values));
318 if !as_node.is_empty() {
319 params.insert("as_node".into(), json!(as_node));
320 }
321 self.send_envelope(new_request_with_id("loop_state_update", params, rid))
322 .await
323 }
324
325 pub async fn send_loop_history_fetch(&self, loop_id: &str, request_id: &[&str]) -> Result<()> {
327 let rid = opt_request_id(request_id);
328 let mut params = Map::new();
329 params.insert("loop_id".into(), json!(loop_id));
330 self.send_envelope(new_request_with_id("loop_history_fetch", params, rid))
331 .await
332 }
333
334 pub async fn send_loop_execution_state_fetch(
337 &self,
338 loop_id: &str,
339 request_id: &[&str],
340 ) -> Result<()> {
341 let rid = opt_request_id(request_id);
342 let mut params = Map::new();
343 params.insert("loop_id".into(), json!(loop_id));
344 self.send_envelope(new_request_with_id(
345 "loop_execution_state_fetch",
346 params,
347 rid,
348 ))
349 .await
350 }
351
352 pub async fn send_auth(
354 &self,
355 access_key: &str,
356 secret_key: &str,
357 request_id: &[&str],
358 ) -> Result<()> {
359 let rid = opt_request_id(request_id);
360 let mut params = Map::new();
361 params.insert("access_key".into(), json!(access_key));
362 params.insert("secret_key".into(), json!(secret_key));
363 self.send_envelope(new_request_with_id("auth", params, rid))
364 .await
365 }
366
367 pub async fn send_auth_refresh(&self, refresh_token: &str, request_id: &[&str]) -> Result<()> {
369 let rid = opt_request_id(request_id);
370 let mut params = Map::new();
371 params.insert("refresh_token".into(), json!(refresh_token));
372 self.send_envelope(new_request_with_id("auth_refresh", params, rid))
373 .await
374 }
375
376 pub async fn send_cron_add(
378 &self,
379 text: &str,
380 priority: i32,
381 request_id: &[&str],
382 ) -> Result<()> {
383 let rid = opt_request_id(request_id);
384 let mut params = Map::new();
385 params.insert("text".into(), json!(text));
386 if priority > 0 {
387 params.insert("priority".into(), json!(priority));
388 }
389 self.send_envelope(new_request_with_id("cron_add", params, rid))
390 .await
391 }
392
393 pub async fn send_cron_list(&self, status: &str, request_id: &[&str]) -> Result<()> {
395 let rid = opt_request_id(request_id);
396 let mut params = Map::new();
397 if !status.is_empty() {
398 params.insert("status".into(), json!(status));
399 }
400 self.send_envelope(new_request_with_id("cron_list", params, rid))
401 .await
402 }
403
404 pub async fn send_cron_show(&self, job_id: &str, request_id: &[&str]) -> Result<()> {
406 let rid = opt_request_id(request_id);
407 let mut params = Map::new();
408 params.insert("job_id".into(), json!(job_id));
409 self.send_envelope(new_request_with_id("cron_show", params, rid))
410 .await
411 }
412
413 pub async fn send_cron_cancel(&self, job_id: &str, request_id: &[&str]) -> Result<()> {
415 let rid = opt_request_id(request_id);
416 let mut params = Map::new();
417 params.insert("job_id".into(), json!(job_id));
418 self.send_envelope(new_request_with_id("cron_cancel", params, rid))
419 .await
420 }
421
422 pub async fn send_job_create(
424 &self,
425 goal: &str,
426 workspace: Option<&str>,
427 request_id: &[&str],
428 ) -> Result<()> {
429 if goal.is_empty() {
430 return Err(crate::errors::Error::msg("goal is required"));
431 }
432 let rid = opt_request_id(request_id);
433 let mut params = Map::new();
434 params.insert("goal".into(), json!(goal));
435 if let Some(ws) = workspace {
436 if !ws.is_empty() {
437 params.insert("workspace".into(), json!(ws));
438 }
439 }
440 self.send_envelope(new_request_with_id("job_create", params, rid))
441 .await
442 }
443
444 pub async fn send_job_status(&self, job_id: &str, request_id: &[&str]) -> Result<()> {
446 if job_id.is_empty() {
447 return Err(crate::errors::Error::msg("job_id is required"));
448 }
449 let rid = opt_request_id(request_id);
450 let mut params = Map::new();
451 params.insert("job_id".into(), json!(job_id));
452 self.send_envelope(new_request_with_id("job_status", params, rid))
453 .await
454 }
455
456 pub async fn send_job_pause(&self, job_id: &str, request_id: &[&str]) -> Result<()> {
458 if job_id.is_empty() {
459 return Err(crate::errors::Error::msg("job_id is required"));
460 }
461 let rid = opt_request_id(request_id);
462 let mut params = Map::new();
463 params.insert("job_id".into(), json!(job_id));
464 self.send_envelope(new_request_with_id("job_pause", params, rid))
465 .await
466 }
467
468 pub async fn send_job_resume(&self, job_id: &str, request_id: &[&str]) -> Result<()> {
470 if job_id.is_empty() {
471 return Err(crate::errors::Error::msg("job_id is required"));
472 }
473 let rid = opt_request_id(request_id);
474 let mut params = Map::new();
475 params.insert("job_id".into(), json!(job_id));
476 self.send_envelope(new_request_with_id("job_resume", params, rid))
477 .await
478 }
479
480 pub async fn send_job_cancel(&self, job_id: &str, request_id: &[&str]) -> Result<()> {
482 if job_id.is_empty() {
483 return Err(crate::errors::Error::msg("job_id is required"));
484 }
485 let rid = opt_request_id(request_id);
486 let mut params = Map::new();
487 params.insert("job_id".into(), json!(job_id));
488 self.send_envelope(new_request_with_id("job_cancel", params, rid))
489 .await
490 }
491
492 pub async fn send_job_dag(&self, job_id: &str, request_id: &[&str]) -> Result<()> {
494 if job_id.is_empty() {
495 return Err(crate::errors::Error::msg("job_id is required"));
496 }
497 let rid = opt_request_id(request_id);
498 let mut params = Map::new();
499 params.insert("job_id".into(), json!(job_id));
500 self.send_envelope(new_request_with_id("job_dag", params, rid))
501 .await
502 }
503}
504
505#[cfg(test)]
506mod tests {
507 use super::*;
508
509 #[test]
510 fn opt_request_id_picks_first_non_empty() {
511 assert_eq!(opt_request_id(&[]), "");
512 assert_eq!(opt_request_id(&[""]), "");
513 assert_eq!(opt_request_id(&["abc"]), "abc");
514 assert_eq!(opt_request_id(&["", "def"]), "def");
515 }
516}