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_auth(
336 &self,
337 access_key: &str,
338 secret_key: &str,
339 request_id: &[&str],
340 ) -> Result<()> {
341 let rid = opt_request_id(request_id);
342 let mut params = Map::new();
343 params.insert("access_key".into(), json!(access_key));
344 params.insert("secret_key".into(), json!(secret_key));
345 self.send_envelope(new_request_with_id("auth", params, rid))
346 .await
347 }
348
349 pub async fn send_auth_refresh(&self, refresh_token: &str, request_id: &[&str]) -> Result<()> {
351 let rid = opt_request_id(request_id);
352 let mut params = Map::new();
353 params.insert("refresh_token".into(), json!(refresh_token));
354 self.send_envelope(new_request_with_id("auth_refresh", params, rid))
355 .await
356 }
357
358 pub async fn send_cron_add(
360 &self,
361 text: &str,
362 priority: i32,
363 request_id: &[&str],
364 ) -> Result<()> {
365 let rid = opt_request_id(request_id);
366 let mut params = Map::new();
367 params.insert("text".into(), json!(text));
368 if priority > 0 {
369 params.insert("priority".into(), json!(priority));
370 }
371 self.send_envelope(new_request_with_id("cron_add", params, rid))
372 .await
373 }
374
375 pub async fn send_cron_list(&self, status: &str, request_id: &[&str]) -> Result<()> {
377 let rid = opt_request_id(request_id);
378 let mut params = Map::new();
379 if !status.is_empty() {
380 params.insert("status".into(), json!(status));
381 }
382 self.send_envelope(new_request_with_id("cron_list", params, rid))
383 .await
384 }
385
386 pub async fn send_cron_show(&self, job_id: &str, request_id: &[&str]) -> Result<()> {
388 let rid = opt_request_id(request_id);
389 let mut params = Map::new();
390 params.insert("job_id".into(), json!(job_id));
391 self.send_envelope(new_request_with_id("cron_show", params, rid))
392 .await
393 }
394
395 pub async fn send_cron_cancel(&self, job_id: &str, request_id: &[&str]) -> Result<()> {
397 let rid = opt_request_id(request_id);
398 let mut params = Map::new();
399 params.insert("job_id".into(), json!(job_id));
400 self.send_envelope(new_request_with_id("cron_cancel", params, rid))
401 .await
402 }
403
404 pub async fn send_job_create(
406 &self,
407 goal: &str,
408 workspace: Option<&str>,
409 request_id: &[&str],
410 ) -> Result<()> {
411 if goal.is_empty() {
412 return Err(crate::errors::Error::msg("goal is required"));
413 }
414 let rid = opt_request_id(request_id);
415 let mut params = Map::new();
416 params.insert("goal".into(), json!(goal));
417 if let Some(ws) = workspace {
418 if !ws.is_empty() {
419 params.insert("workspace".into(), json!(ws));
420 }
421 }
422 self.send_envelope(new_request_with_id("job_create", params, rid))
423 .await
424 }
425
426 pub async fn send_job_status(&self, job_id: &str, request_id: &[&str]) -> Result<()> {
428 if job_id.is_empty() {
429 return Err(crate::errors::Error::msg("job_id is required"));
430 }
431 let rid = opt_request_id(request_id);
432 let mut params = Map::new();
433 params.insert("job_id".into(), json!(job_id));
434 self.send_envelope(new_request_with_id("job_status", params, rid))
435 .await
436 }
437
438 pub async fn send_job_pause(&self, job_id: &str, request_id: &[&str]) -> Result<()> {
440 if job_id.is_empty() {
441 return Err(crate::errors::Error::msg("job_id is required"));
442 }
443 let rid = opt_request_id(request_id);
444 let mut params = Map::new();
445 params.insert("job_id".into(), json!(job_id));
446 self.send_envelope(new_request_with_id("job_pause", params, rid))
447 .await
448 }
449
450 pub async fn send_job_resume(&self, job_id: &str, request_id: &[&str]) -> Result<()> {
452 if job_id.is_empty() {
453 return Err(crate::errors::Error::msg("job_id is required"));
454 }
455 let rid = opt_request_id(request_id);
456 let mut params = Map::new();
457 params.insert("job_id".into(), json!(job_id));
458 self.send_envelope(new_request_with_id("job_resume", params, rid))
459 .await
460 }
461
462 pub async fn send_job_cancel(&self, job_id: &str, request_id: &[&str]) -> Result<()> {
464 if job_id.is_empty() {
465 return Err(crate::errors::Error::msg("job_id is required"));
466 }
467 let rid = opt_request_id(request_id);
468 let mut params = Map::new();
469 params.insert("job_id".into(), json!(job_id));
470 self.send_envelope(new_request_with_id("job_cancel", params, rid))
471 .await
472 }
473
474 pub async fn send_job_dag(&self, job_id: &str, request_id: &[&str]) -> Result<()> {
476 if job_id.is_empty() {
477 return Err(crate::errors::Error::msg("job_id is required"));
478 }
479 let rid = opt_request_id(request_id);
480 let mut params = Map::new();
481 params.insert("job_id".into(), json!(job_id));
482 self.send_envelope(new_request_with_id("job_dag", params, rid))
483 .await
484 }
485}
486
487#[cfg(test)]
488mod tests {
489 use super::*;
490
491 #[test]
492 fn opt_request_id_picks_first_non_empty() {
493 assert_eq!(opt_request_id(&[]), "");
494 assert_eq!(opt_request_id(&[""]), "");
495 assert_eq!(opt_request_id(&["abc"]), "abc");
496 assert_eq!(opt_request_id(&["", "def"]), "def");
497 }
498}