1use std::path::PathBuf;
12use std::time::Duration;
13
14use colored::Colorize;
15
16const REQUEST_TIMEOUT: Duration = Duration::from_secs(15);
17
18#[derive(Debug, Clone, Default)]
20pub struct ConnArgs {
21 pub server_url: Option<String>,
23 pub port: Option<u16>,
25 pub data_dir: Option<PathBuf>,
27}
28
29impl ConnArgs {
30 pub(crate) fn api_base(&self) -> String {
32 if let Some(url) = &self.server_url {
33 let url = url.trim_end_matches('/');
34 let url = if url.contains("://") {
36 url.to_string()
37 } else {
38 format!("http://{url}")
39 };
40 return format!("{url}/api/v1");
41 }
42 let config = bamboo_llm::Config::from_data_dir(self.data_dir.clone());
43 let port = self.port.unwrap_or(config.server.port);
44 let host = match config.server.bind.trim() {
45 "" | "0.0.0.0" | "::" | "[::]" => "127.0.0.1".to_string(),
47 h if h.contains(':') && !h.starts_with('[') => format!("[{h}]"),
49 h => h.to_string(),
50 };
51 format!("http://{host}:{port}/api/v1")
52 }
53}
54
55pub(crate) fn unreachable(base: &str, e: reqwest::Error) -> anyhow::Error {
56 anyhow::anyhow!("could not reach the server at {base} ({e}). Is `bamboo serve` running?")
57}
58
59pub(crate) fn guard_id_segment(kind: &str, id: &str) -> anyhow::Result<()> {
64 if id.is_empty()
65 || id == "."
66 || id == ".."
67 || id.contains(['/', '\\', '?', '#', '%'])
68 || id.chars().any(char::is_whitespace)
69 {
70 anyhow::bail!("invalid {kind}: '{id}'");
71 }
72 Ok(())
73}
74
75pub async fn health(conn: ConnArgs) -> anyhow::Result<()> {
78 let base = conn.api_base();
79 let url = format!("{base}/health");
80 let resp = reqwest::Client::new()
81 .get(&url)
82 .timeout(REQUEST_TIMEOUT)
83 .send()
84 .await
85 .map_err(|e| unreachable(&base, e))?;
86 if resp.status().is_success() {
87 println!("{} {base}", "● healthy".green().bold());
88 Ok(())
89 } else {
90 anyhow::bail!("unhealthy: HTTP {} from {url}", resp.status());
91 }
92}
93
94pub async fn status(conn: ConnArgs) -> anyhow::Result<()> {
96 let base = conn.api_base();
97 let server = base.trim_end_matches("/api/v1");
98 println!("{:<10}{server}", "server:".bold());
99
100 let client = reqwest::Client::new();
101 let health = client
102 .get(format!("{base}/health"))
103 .timeout(REQUEST_TIMEOUT)
104 .send()
105 .await;
106 match health {
107 Ok(r) if r.status().is_success() => println!("{:<10}{}", "health:".bold(), "ok".green()),
108 Ok(r) => {
109 println!(
110 "{:<10}{} (HTTP {})",
111 "health:".bold(),
112 "down".red(),
113 r.status()
114 );
115 return Ok(());
116 }
117 Err(e) => {
118 println!("{:<10}{} ({e})", "health:".bold(), "unreachable".red());
119 return Ok(());
120 }
121 }
122
123 if let Ok(r) = client
124 .get(format!("{base}/sessions"))
125 .timeout(REQUEST_TIMEOUT)
126 .send()
127 .await
128 {
129 if let Ok(v) = r.json::<serde_json::Value>().await {
130 let sessions = v.get("sessions").and_then(|s| s.as_array());
131 let total = sessions.map(|s| s.len()).unwrap_or(0);
132 let running = sessions.map(|s| count_running(s)).unwrap_or(0);
133 println!(
134 "{:<10}{total} total, {} running",
135 "sessions:".bold(),
136 running.to_string().cyan()
137 );
138 }
139 }
140 Ok(())
141}
142
143pub async fn sessions_list(conn: ConnArgs) -> anyhow::Result<()> {
145 let base = conn.api_base();
146 let url = format!("{base}/sessions");
147 let resp = reqwest::Client::new()
148 .get(&url)
149 .timeout(REQUEST_TIMEOUT)
150 .send()
151 .await
152 .map_err(|e| unreachable(&base, e))?;
153 if !resp.status().is_success() {
154 anyhow::bail!("GET {url} -> HTTP {}", resp.status());
155 }
156 let v: serde_json::Value = resp.json().await?;
157 let sessions = v.get("sessions").and_then(|s| s.as_array());
158 let sessions = match sessions {
159 Some(s) if !s.is_empty() => s,
160 _ => {
161 println!("(no sessions)");
162 return Ok(());
163 }
164 };
165
166 println!(
169 "{:<38} {:<5} {:<26} {:>5} TITLE",
170 "SESSION ID", "RUN", "MODEL", "MSGS"
171 );
172 for s in sessions {
173 let id = s.get("id").and_then(|x| x.as_str()).unwrap_or("?");
174 let running = s
175 .get("is_running")
176 .and_then(|b| b.as_bool())
177 .unwrap_or(false);
178 let model = s.get("model").and_then(|x| x.as_str()).unwrap_or("");
179 let msgs = s.get("message_count").and_then(|x| x.as_u64()).unwrap_or(0);
180 let title = s.get("title").and_then(|x| x.as_str()).unwrap_or("");
181 println!(
182 "{:<38} {:<5} {:<26} {:>5} {}",
183 id,
184 if running { "run" } else { "-" },
185 truncate(model, 26),
186 msgs,
187 truncate(title, 60)
188 );
189 }
190 let running = count_running(sessions);
191 println!(
192 "\n{running} running. Stop one with: {}",
193 "bamboo stop <session-id>".cyan()
194 );
195 Ok(())
196}
197
198pub async fn stop(conn: ConnArgs, session_id: &str) -> anyhow::Result<()> {
200 guard_id_segment("session id", session_id)?;
201 let base = conn.api_base();
202 let url = format!("{base}/stop/{session_id}");
203 let resp = reqwest::Client::new()
204 .post(&url)
205 .timeout(REQUEST_TIMEOUT)
206 .send()
207 .await
208 .map_err(|e| unreachable(&base, e))?;
209 let status = resp.status();
210 let body: serde_json::Value = resp.json().await.unwrap_or(serde_json::Value::Null);
211 let message = body
212 .get("message")
213 .and_then(|m| m.as_str())
214 .unwrap_or("")
215 .to_string();
216 if status.is_success() {
217 let msg = if message.is_empty() {
218 "stopped"
219 } else {
220 &message
221 };
222 println!("{} {msg}", "✓".green());
223 Ok(())
224 } else if status.as_u16() == 404 {
225 anyhow::bail!(
226 "no active run for session '{session_id}'{}",
227 if message.is_empty() {
228 String::new()
229 } else {
230 format!(" ({message})")
231 }
232 );
233 } else {
234 anyhow::bail!("stop failed: HTTP {status} {message}");
235 }
236}
237
238pub async fn history(conn: ConnArgs, session_id: &str) -> anyhow::Result<()> {
243 guard_id_segment("session id", session_id)?;
244 let base = conn.api_base();
245 let url = format!("{base}/history/{session_id}");
246 let resp = reqwest::Client::new()
247 .get(&url)
248 .timeout(REQUEST_TIMEOUT)
249 .send()
250 .await
251 .map_err(|e| unreachable(&base, e))?;
252 if resp.status().as_u16() == 404 {
253 anyhow::bail!("session '{session_id}' not found");
254 }
255 if !resp.status().is_success() {
256 anyhow::bail!("GET {url} -> HTTP {}", resp.status());
257 }
258 let v: serde_json::Value = resp.json().await?;
259 let messages = v.get("messages").and_then(|m| m.as_array());
260 let messages = match messages {
261 Some(m) if !m.is_empty() => m,
262 _ => {
263 println!("(no messages)");
264 return Ok(());
265 }
266 };
267 for m in messages {
268 let role = m.get("role").and_then(|r| r.as_str()).unwrap_or("?");
269 let content = m.get("content").and_then(|c| c.as_str()).unwrap_or("");
270 let label = match role {
271 "user" => "user".cyan(),
272 "assistant" => "assistant".green(),
273 "system" => "system".dimmed(),
274 "tool" => "tool".yellow(),
275 other => other.normal(),
276 };
277 println!("{label}: {content}");
278 }
279 let total = v
283 .get("total_message_count")
284 .and_then(|x| x.as_u64())
285 .unwrap_or(messages.len() as u64);
286 let truncated = v
287 .get("truncated")
288 .and_then(|x| x.as_bool())
289 .unwrap_or(false);
290 println!(
291 "\n{}",
292 history_summary(session_id, messages.len(), total, truncated)
293 );
294 Ok(())
295}
296
297fn history_summary(session_id: &str, shown: usize, total: u64, truncated: bool) -> String {
301 if truncated {
302 format!(
303 "{total} message(s) in session {session_id} (showing the newest {shown}; older messages omitted by the server's history cap)."
304 )
305 } else {
306 format!("{total} message(s) in session {session_id}.")
307 }
308}
309
310pub async fn respond(conn: ConnArgs, session_id: &str, answer: &str) -> anyhow::Result<()> {
314 guard_id_segment("session id", session_id)?;
315 let base = conn.api_base();
316 let url = format!("{base}/respond/{session_id}");
317 let resp = reqwest::Client::new()
318 .post(&url)
319 .timeout(REQUEST_TIMEOUT)
320 .json(&serde_json::json!({ "response": answer }))
323 .send()
324 .await
325 .map_err(|e| unreachable(&base, e))?;
326 let status = resp.status();
327 let body: serde_json::Value = resp.json().await.unwrap_or(serde_json::Value::Null);
328 if status.is_success() {
329 let auto_resume = body
330 .get("auto_resume_status")
331 .and_then(|s| s.as_str())
332 .unwrap_or("unknown");
333 println!(
334 "{} response recorded; the run resumes server-side (auto-resume: {auto_resume}).",
335 "✓".green()
336 );
337 if let Some(run_id) = body.get("run_id").and_then(|r| r.as_str()) {
338 println!("run id: {run_id}");
339 }
340 return Ok(());
341 }
342 if status.as_u16() == 404 {
343 anyhow::bail!("session '{session_id}' not found");
344 }
345 let error = server_error_message(&body);
346 if status.as_u16() == 400 && error.contains("No pending question") {
347 anyhow::bail!(
348 "session '{session_id}' has no pending question — nothing to answer \
349 (check with: bamboo respond {session_id} --pending)"
350 );
351 }
352 let detail = body.get("message").and_then(|m| m.as_str()).unwrap_or("");
353 anyhow::bail!(
354 "respond failed: HTTP {status}{}{}",
355 if error.is_empty() {
356 String::new()
357 } else {
358 format!(" {error}")
359 },
360 if detail.is_empty() {
361 String::new()
362 } else {
363 format!(" ({detail})")
364 }
365 );
366}
367
368pub async fn respond_pending(conn: ConnArgs, session_id: &str, json: bool) -> anyhow::Result<()> {
371 guard_id_segment("session id", session_id)?;
372 let base = conn.api_base();
373 let url = format!("{base}/respond/{session_id}/pending");
374 let resp = reqwest::Client::new()
375 .get(&url)
376 .timeout(REQUEST_TIMEOUT)
377 .send()
378 .await
379 .map_err(|e| unreachable(&base, e))?;
380 if resp.status().as_u16() == 404 {
381 anyhow::bail!("session '{session_id}' not found");
382 }
383 if !resp.status().is_success() {
384 anyhow::bail!("GET {url} -> HTTP {}", resp.status());
385 }
386 let v: serde_json::Value = resp.json().await?;
387 if json {
388 println!("{}", serde_json::to_string_pretty(&v)?);
389 return Ok(());
390 }
391 match format_pending_question(session_id, &v) {
392 Some(text) => println!("{text}"),
393 None => println!("no pending question for session '{session_id}'."),
394 }
395 Ok(())
396}
397
398fn format_pending_question(session_id: &str, v: &serde_json::Value) -> Option<String> {
401 if !v
402 .get("has_pending_question")
403 .and_then(|b| b.as_bool())
404 .unwrap_or(false)
405 {
406 return None;
407 }
408 let question = v.get("question").and_then(|q| q.as_str()).unwrap_or("");
409 let mut out = format!("session: {session_id}\nquestion: {question}\n");
410 if let Some(options) = v.get("options").and_then(|o| o.as_array()) {
411 if !options.is_empty() {
412 out.push_str("options:\n");
413 for (i, opt) in options.iter().enumerate() {
414 let opt = opt
415 .as_str()
416 .map(str::to_string)
417 .unwrap_or_else(|| opt.to_string());
418 out.push_str(&format!(" {}. {opt}\n", i + 1));
419 }
420 }
421 }
422 if v.get("allow_custom")
423 .and_then(|b| b.as_bool())
424 .unwrap_or(false)
425 {
426 out.push_str("(custom free-text answers are allowed)\n");
427 }
428 if let Some(tool) = v
429 .get("tool_name")
430 .and_then(|t| t.as_str())
431 .filter(|t| !t.is_empty())
432 {
433 out.push_str(&format!("tool: {tool}\n"));
434 }
435 out.push_str(&format!(
436 "\nAnswer with: bamboo respond {session_id} \"<answer>\" — answering resumes the run server-side."
437 ));
438 Some(out)
439}
440
441pub async fn session_show(conn: ConnArgs, session_id: &str, json: bool) -> anyhow::Result<()> {
444 guard_id_segment("session id", session_id)?;
445 let base = conn.api_base();
446 let url = format!("{base}/sessions/{session_id}");
447 let resp = reqwest::Client::new()
448 .get(&url)
449 .timeout(REQUEST_TIMEOUT)
450 .send()
451 .await
452 .map_err(|e| unreachable(&base, e))?;
453 if resp.status().as_u16() == 404 {
454 anyhow::bail!("session '{session_id}' not found");
455 }
456 if !resp.status().is_success() {
457 anyhow::bail!("GET {url} -> HTTP {}", resp.status());
458 }
459 let v: serde_json::Value = resp.json().await?;
460 if json {
461 println!("{}", serde_json::to_string_pretty(&v)?);
462 return Ok(());
463 }
464 let session = v.get("session").unwrap_or(&v);
465 println!("{}", format_session_detail(session));
466 Ok(())
467}
468
469fn format_session_detail(s: &serde_json::Value) -> String {
473 let str_field = |key: &str| s.get(key).and_then(|x| x.as_str()).unwrap_or("");
474 let mut lines: Vec<String> = Vec::new();
475 let mut push = |label: &str, value: String| {
476 let label = format!("{:<16}", format!("{label}:"));
479 lines.push(format!("{}{value}", label.bold()));
480 };
481
482 push("id", str_field("id").to_string());
483 push("title", str_field("title").to_string());
484 push("kind", str_field("kind").to_string());
485 let model = str_field("model").to_string();
486 let model = match s.get("provider").and_then(|p| p.as_str()) {
487 Some(provider) if !provider.is_empty() => format!("{provider}:{model}"),
488 _ => model,
489 };
490 push("model", model);
491 let running = s
492 .get("is_running")
493 .and_then(|b| b.as_bool())
494 .unwrap_or(false);
495 push(
496 "running",
497 if running {
498 "yes".green().to_string()
499 } else {
500 "no".to_string()
501 },
502 );
503 if let Some(status) = s.get("last_run_status").and_then(|x| x.as_str()) {
504 let mut line = status.to_string();
505 if let Some(err) = s.get("last_run_error").and_then(|x| x.as_str()) {
506 line.push_str(&format!(" ({err})"));
507 }
508 push("last run", line);
509 }
510 let pending = s
511 .get("has_pending_question")
512 .and_then(|b| b.as_bool())
513 .unwrap_or(false);
514 push(
515 "pending q",
516 if pending {
517 "yes (see: bamboo respond <id> --pending)"
518 .yellow()
519 .to_string()
520 } else {
521 "no".to_string()
522 },
523 );
524 push(
525 "messages",
526 s.get("message_count")
527 .and_then(|x| x.as_u64())
528 .unwrap_or(0)
529 .to_string(),
530 );
531 if s.get("pinned").and_then(|b| b.as_bool()).unwrap_or(false) {
532 push("pinned", "yes".to_string());
533 }
534 if let Some(parent) = s.get("parent_session_id").and_then(|x| x.as_str()) {
535 push("parent", parent.to_string());
536 }
537 let child_count = s
538 .get("running_child_count")
539 .and_then(|x| x.as_u64())
540 .unwrap_or(0);
541 if child_count > 0 {
542 push("children", format!("{child_count} running"));
543 }
544 push("created", str_field("created_at").to_string());
545 push("last activity", str_field("last_activity_at").to_string());
546 if let Some(placement) = s.get("placement") {
547 let kind = placement.get("kind").and_then(|x| x.as_str()).unwrap_or("");
548 let host = placement.get("host").and_then(|x| x.as_str()).unwrap_or("");
549 if !kind.is_empty() || !host.is_empty() {
550 push("placement", format!("{kind} @ {host}"));
551 }
552 }
553 lines.join("\n")
554}
555
556pub async fn session_delete(conn: ConnArgs, session_id: &str, yes: bool) -> anyhow::Result<()> {
560 guard_id_segment("session id", session_id)?;
561 if !yes && !confirm(&format!(
562 "Delete session '{session_id}'? This cancels any running execution and removes it permanently."
563 ))? {
564 println!("aborted (nothing deleted).");
565 return Ok(());
566 }
567 let base = conn.api_base();
568 let url = format!("{base}/sessions/{session_id}");
569 let resp = reqwest::Client::new()
570 .delete(&url)
571 .timeout(REQUEST_TIMEOUT)
572 .send()
573 .await
574 .map_err(|e| unreachable(&base, e))?;
575 let status = resp.status();
576 if status.is_success() {
577 println!("{} session '{session_id}' deleted", "✓".green());
578 return Ok(());
579 }
580 if status.as_u16() == 404 {
581 anyhow::bail!("session '{session_id}' not found");
582 }
583 let body: serde_json::Value = resp.json().await.unwrap_or(serde_json::Value::Null);
584 let error = server_error_message(&body);
585 anyhow::bail!(
586 "delete failed: HTTP {status}{}",
587 if error.is_empty() {
588 String::new()
589 } else {
590 format!(" {error}")
591 }
592 );
593}
594
595pub(crate) fn confirm(prompt: &str) -> anyhow::Result<bool> {
599 use std::io::Write as _;
600 print!("{prompt} [y/N] ");
601 std::io::stdout().flush()?;
602 let mut line = String::new();
603 std::io::stdin().read_line(&mut line)?;
604 let answer = line.trim().to_ascii_lowercase();
605 Ok(answer == "y" || answer == "yes")
606}
607
608#[derive(Debug, Clone, Default)]
618pub struct ScheduleCreateArgs {
619 pub name: Option<String>,
621 pub cron: Option<String>,
623 pub every: Option<u64>,
625 pub daily: Option<String>,
627 pub prompt: Option<String>,
629 pub model: Option<String>,
631 pub workspace: Option<String>,
633 pub timezone: Option<String>,
635 pub disabled: bool,
637 pub json: Option<String>,
640}
641
642pub async fn schedules_list(conn: ConnArgs, json: bool) -> anyhow::Result<()> {
644 let base = conn.api_base();
645 let body = get_json(&base, &format!("{base}/schedules")).await?;
646 if json {
647 println!("{}", serde_json::to_string_pretty(&body)?);
648 return Ok(());
649 }
650 let schedules = body.get("schedules").and_then(|s| s.as_array());
651 let schedules = match schedules {
652 Some(s) if !s.is_empty() => s,
653 _ => {
654 println!("(no schedules)");
655 return Ok(());
656 }
657 };
658
659 println!(
661 "{:<38} {:<4} {:<24} {:<20} {:<20} NAME",
662 "SCHEDULE ID", "ON", "TRIGGER", "NEXT RUN", "LAST RUN"
663 );
664 for s in schedules {
665 let id = s.get("id").and_then(|x| x.as_str()).unwrap_or("?");
666 let enabled = s.get("enabled").and_then(|b| b.as_bool()).unwrap_or(false);
667 let trigger = s.get("trigger").map(trigger_summary).unwrap_or_default();
668 let state = s.get("state");
669 let next = state.and_then(|st| st.get("next_fire_at"));
670 let last = state.and_then(|st| st.get("last_started_at"));
671 let name = s.get("name").and_then(|x| x.as_str()).unwrap_or("");
672 println!(
673 "{:<38} {:<4} {:<24} {:<20} {:<20} {}",
674 id,
675 if enabled { "on" } else { "off" },
676 truncate(&trigger, 24),
677 fmt_ts(next),
678 fmt_ts(last),
679 truncate(name, 40)
680 );
681 }
682 println!(
683 "\n{} schedule(s). Inspect one with: {}",
684 schedules.len(),
685 "bamboo schedules show <id>".cyan()
686 );
687 Ok(())
688}
689
690pub async fn schedules_show(conn: ConnArgs, schedule_id: &str, json: bool) -> anyhow::Result<()> {
693 guard_id_segment("schedule id", schedule_id)?;
694 let base = conn.api_base();
695 let schedule = find_schedule(&base, schedule_id).await?;
696 if json {
697 println!("{}", serde_json::to_string_pretty(&schedule)?);
698 return Ok(());
699 }
700
701 let str_of = |v: &serde_json::Value| v.as_str().map(str::to_string);
702 let field = |key: &str| schedule.get(key).and_then(str_of).unwrap_or_default();
703 let enabled = schedule
704 .get("enabled")
705 .and_then(|b| b.as_bool())
706 .unwrap_or(false);
707 println!("{:<16}{}", "id:".bold(), field("id"));
708 println!("{:<16}{}", "name:".bold(), field("name"));
709 println!(
710 "{:<16}{}",
711 "enabled:".bold(),
712 if enabled {
713 "true".green()
714 } else {
715 "false".red()
716 }
717 );
718 if let Some(trigger) = schedule.get("trigger") {
719 println!("{:<16}{}", "trigger:".bold(), trigger_summary(trigger));
720 }
721 for key in ["timezone", "start_at", "end_at"] {
722 if let Some(value) = schedule.get(key).and_then(|v| v.as_str()) {
723 println!("{:<16}{value}", format!("{key}:").bold());
724 }
725 }
726 for key in ["misfire_policy", "overlap_policy"] {
727 if let Some(value) = schedule.get(key) {
728 let rendered = value
729 .get("type")
730 .and_then(|t| t.as_str())
731 .map(str::to_string)
732 .or_else(|| str_of(value))
733 .unwrap_or_else(|| value.to_string());
734 println!("{:<16}{rendered}", format!("{key}:").bold());
735 }
736 }
737 if let Some(state) = schedule.get("state") {
738 println!(
739 "{:<16}{}",
740 "next fire:".bold(),
741 fmt_ts(state.get("next_fire_at"))
742 );
743 println!(
744 "{:<16}{}",
745 "last started:".bold(),
746 fmt_ts(state.get("last_started_at"))
747 );
748 println!(
749 "{:<16}{}",
750 "last success:".bold(),
751 fmt_ts(state.get("last_success_at"))
752 );
753 let count = |key: &str| state.get(key).and_then(|v| v.as_u64()).unwrap_or(0);
754 println!(
755 "{:<16}{} total, {} ok, {} failed, {} missed ({} queued, {} running now)",
756 "runs:".bold(),
757 count("total_run_count"),
758 count("total_success_count"),
759 count("total_failure_count"),
760 count("total_missed_count"),
761 count("queued_run_count"),
762 count("running_run_count"),
763 );
764 }
765 if let Some(rc) = schedule.get("run_config") {
766 let rc_str = |key: &str| rc.get(key).and_then(|v| v.as_str());
767 if let Some(task) = rc_str("task_message") {
768 println!("{:<16}{}", "prompt:".bold(), truncate(task, 120));
769 }
770 for (label, key) in [
771 ("model:", "model"),
772 ("workspace:", "workspace_path"),
773 ("reasoning:", "reasoning_effort"),
774 ] {
775 if let Some(value) = rc_str(key) {
776 println!("{:<16}{value}", label.bold());
777 }
778 }
779 let auto = rc
780 .get("auto_execute")
781 .and_then(|b| b.as_bool())
782 .unwrap_or(false);
783 println!("{:<16}{auto}", "auto-execute:".bold());
784 }
785 println!(
786 "{:<16}{} {:<10}{}",
787 "created:".bold(),
788 fmt_ts(schedule.get("created_at")),
789 "updated:".bold(),
790 fmt_ts(schedule.get("updated_at"))
791 );
792 Ok(())
793}
794
795pub async fn schedules_create(conn: ConnArgs, args: ScheduleCreateArgs) -> anyhow::Result<()> {
799 let payload = match &args.json {
800 Some(source) => read_json_payload(source)?,
801 None => build_create_payload(&args)?,
802 };
803
804 let base = conn.api_base();
805 let url = format!("{base}/schedules");
806 let resp = reqwest::Client::new()
807 .post(&url)
808 .timeout(REQUEST_TIMEOUT)
809 .json(&payload)
810 .send()
811 .await
812 .map_err(|e| unreachable(&base, e))?;
813 let status = resp.status();
814 let body: serde_json::Value = resp.json().await.unwrap_or(serde_json::Value::Null);
815 if !status.is_success() {
816 anyhow::bail!(
817 "create failed: HTTP {status} {}",
818 server_error_message(&body)
819 );
820 }
821 let id = body.get("id").and_then(|x| x.as_str()).unwrap_or("?");
822 let name = body.get("name").and_then(|x| x.as_str()).unwrap_or("");
823 let enabled = body
824 .get("enabled")
825 .and_then(|b| b.as_bool())
826 .unwrap_or(false);
827 let trigger = body.get("trigger").map(trigger_summary).unwrap_or_default();
828 println!(
829 "{} created schedule {id} ('{name}', {trigger}, {})",
830 "✓".green(),
831 if enabled { "enabled" } else { "disabled" }
832 );
833 if let Some(next) = body.get("state").and_then(|st| st.get("next_fire_at")) {
834 println!(" next fire: {}", fmt_ts(Some(next)));
835 }
836 Ok(())
837}
838
839pub async fn schedules_delete(conn: ConnArgs, schedule_id: &str, yes: bool) -> anyhow::Result<()> {
842 guard_id_segment("schedule id", schedule_id)?;
843 let base = conn.api_base();
844 if !yes {
845 let schedule = find_schedule(&base, schedule_id).await?;
848 let name = schedule.get("name").and_then(|x| x.as_str()).unwrap_or("?");
849 if !confirm(&format!("Delete schedule '{name}' ({schedule_id})?"))? {
850 println!("aborted (nothing deleted).");
851 return Ok(());
852 }
853 }
854 let url = format!("{base}/schedules/{schedule_id}");
855 let resp = reqwest::Client::new()
856 .delete(&url)
857 .timeout(REQUEST_TIMEOUT)
858 .send()
859 .await
860 .map_err(|e| unreachable(&base, e))?;
861 let status = resp.status();
862 let body: serde_json::Value = resp.json().await.unwrap_or(serde_json::Value::Null);
863 if status.as_u16() == 404 {
864 anyhow::bail!("schedule '{schedule_id}' not found");
865 }
866 if !status.is_success() {
867 anyhow::bail!(
868 "delete failed: HTTP {status} {}",
869 server_error_message(&body)
870 );
871 }
872 println!("{} deleted schedule {schedule_id}", "✓".green());
873 Ok(())
874}
875
876pub async fn schedules_run(conn: ConnArgs, schedule_id: &str) -> anyhow::Result<()> {
878 guard_id_segment("schedule id", schedule_id)?;
879 let base = conn.api_base();
880 let url = format!("{base}/schedules/{schedule_id}/run");
881 let resp = reqwest::Client::new()
882 .post(&url)
883 .timeout(REQUEST_TIMEOUT)
884 .send()
885 .await
886 .map_err(|e| unreachable(&base, e))?;
887 let status = resp.status();
888 let body: serde_json::Value = resp.json().await.unwrap_or(serde_json::Value::Null);
889 if status.as_u16() == 404 {
890 anyhow::bail!("schedule '{schedule_id}' not found");
891 }
892 if !status.is_success() {
893 anyhow::bail!("run failed: HTTP {status} {}", server_error_message(&body));
894 }
895 let run_id = body.get("run_id").and_then(|x| x.as_str()).unwrap_or("?");
896 println!(
897 "{} run {run_id} enqueued (watch it with: {})",
898 "✓".green(),
899 format!("bamboo schedules runs {schedule_id}").cyan()
900 );
901 Ok(())
902}
903
904pub async fn schedules_runs(conn: ConnArgs, schedule_id: &str, json: bool) -> anyhow::Result<()> {
906 guard_id_segment("schedule id", schedule_id)?;
907 let base = conn.api_base();
908 let body = get_json(&base, &format!("{base}/schedules/{schedule_id}/runs")).await?;
909 if json {
910 println!("{}", serde_json::to_string_pretty(&body)?);
911 return Ok(());
912 }
913 let runs = body.get("runs").and_then(|r| r.as_array());
914 let runs = match runs {
915 Some(r) if !r.is_empty() => r,
916 _ => {
917 println!("(no runs)");
918 return Ok(());
919 }
920 };
921 println!(
922 "{:<38} {:<9} {:<20} {:<20} {:>9} SESSION",
923 "RUN ID", "STATUS", "SCHEDULED FOR", "STARTED", "DURATION"
924 );
925 for r in runs {
926 let run_id = r.get("run_id").and_then(|x| x.as_str()).unwrap_or("?");
927 let run_status = r.get("status").and_then(|x| x.as_str()).unwrap_or("?");
928 let duration = r
929 .get("execution_duration_ms")
930 .and_then(|x| x.as_u64())
931 .map(|ms| format!("{ms}ms"))
932 .unwrap_or_else(|| "-".to_string());
933 let session = r.get("session_id").and_then(|x| x.as_str()).unwrap_or("-");
934 println!(
935 "{:<38} {:<9} {:<20} {:<20} {:>9} {}",
936 run_id,
937 run_status,
938 fmt_ts(r.get("scheduled_for")),
939 fmt_ts(r.get("started_at")),
940 duration,
941 session
942 );
943 }
944 println!("\n{} run(s) for schedule {schedule_id}.", runs.len());
945 Ok(())
946}
947
948async fn get_json(base: &str, url: &str) -> anyhow::Result<serde_json::Value> {
951 let resp = reqwest::Client::new()
952 .get(url)
953 .timeout(REQUEST_TIMEOUT)
954 .send()
955 .await
956 .map_err(|e| unreachable(base, e))?;
957 let status = resp.status();
958 let body: serde_json::Value = resp.json().await.unwrap_or(serde_json::Value::Null);
959 if !status.is_success() {
960 anyhow::bail!("GET {url} -> HTTP {status} {}", server_error_message(&body));
961 }
962 Ok(body)
963}
964
965async fn find_schedule(base: &str, schedule_id: &str) -> anyhow::Result<serde_json::Value> {
968 let body = get_json(base, &format!("{base}/schedules")).await?;
969 body.get("schedules")
970 .and_then(|s| s.as_array())
971 .and_then(|schedules| {
972 schedules
973 .iter()
974 .find(|s| s.get("id").and_then(|x| x.as_str()) == Some(schedule_id))
975 })
976 .cloned()
977 .ok_or_else(|| {
978 anyhow::anyhow!(
979 "schedule '{schedule_id}' not found (list them with: bamboo schedules list)"
980 )
981 })
982}
983
984fn build_create_payload(args: &ScheduleCreateArgs) -> anyhow::Result<serde_json::Value> {
986 let name = args
987 .name
988 .as_deref()
989 .ok_or_else(|| anyhow::anyhow!("--name is required (or pass --json)"))?;
990 let prompt = args
991 .prompt
992 .as_deref()
993 .ok_or_else(|| anyhow::anyhow!("--prompt is required (or pass --json)"))?;
994
995 let trigger = if let Some(expr) = &args.cron {
996 serde_json::json!({ "type": "cron", "expr": expr })
997 } else if let Some(every_seconds) = args.every {
998 serde_json::json!({ "type": "interval", "every_seconds": every_seconds })
999 } else if let Some(hms) = &args.daily {
1000 let (hour, minute, second) = parse_daily_time(hms)?;
1001 serde_json::json!({ "type": "daily", "hour": hour, "minute": minute, "second": second })
1002 } else {
1003 anyhow::bail!("a trigger is required: --cron <expr>, --every <seconds>, or --daily <HH:MM[:SS]> (or pass --json)");
1004 };
1005
1006 let mut run_config = serde_json::json!({
1009 "task_message": prompt,
1010 "auto_execute": true,
1011 });
1012 if let Some(model) = &args.model {
1013 run_config["model"] = serde_json::json!(model);
1014 }
1015 if let Some(workspace) = &args.workspace {
1016 run_config["workspace_path"] = serde_json::json!(workspace);
1017 }
1018
1019 let mut payload = serde_json::json!({
1020 "name": name,
1021 "trigger": trigger,
1022 "enabled": !args.disabled,
1023 "run_config": run_config,
1024 });
1025 if let Some(timezone) = &args.timezone {
1026 payload["timezone"] = serde_json::json!(timezone);
1027 }
1028 Ok(payload)
1029}
1030
1031fn parse_daily_time(value: &str) -> anyhow::Result<(u8, u8, u8)> {
1033 let bad = || anyhow::anyhow!("invalid --daily time '{value}' (expected HH:MM or HH:MM:SS)");
1034 let parts: Vec<&str> = value.split(':').collect();
1035 if parts.len() != 2 && parts.len() != 3 {
1036 return Err(bad());
1037 }
1038 let hour: u8 = parts[0].parse().map_err(|_| bad())?;
1039 let minute: u8 = parts[1].parse().map_err(|_| bad())?;
1040 let second: u8 = if parts.len() == 3 {
1041 parts[2].parse().map_err(|_| bad())?
1042 } else {
1043 0
1044 };
1045 if hour > 23 || minute > 59 || second > 59 {
1046 return Err(bad());
1047 }
1048 Ok((hour, minute, second))
1049}
1050
1051fn read_json_payload(source: &str) -> anyhow::Result<serde_json::Value> {
1053 let text = if source == "-" {
1054 use std::io::Read as _;
1055 let mut buf = String::new();
1056 std::io::stdin().read_to_string(&mut buf)?;
1057 buf
1058 } else {
1059 std::fs::read_to_string(source)
1060 .map_err(|e| anyhow::anyhow!("failed to read '{source}': {e}"))?
1061 };
1062 serde_json::from_str(text.trim()).map_err(|e| anyhow::anyhow!("payload is not valid JSON: {e}"))
1063}
1064
1065pub(crate) fn server_error_message(body: &serde_json::Value) -> String {
1069 body.get("error")
1070 .and_then(|error| {
1071 error
1072 .get("message")
1073 .and_then(serde_json::Value::as_str)
1074 .or_else(|| error.as_str())
1075 })
1076 .map(|e| format!("({e})"))
1077 .unwrap_or_default()
1078}
1079
1080fn trigger_summary(trigger: &serde_json::Value) -> String {
1082 let joined = |key: &str| {
1083 trigger
1084 .get(key)
1085 .and_then(|v| v.as_array())
1086 .map(|items| {
1087 items
1088 .iter()
1089 .map(|d| match d {
1090 serde_json::Value::String(s) => s.clone(),
1091 other => other.to_string(),
1092 })
1093 .collect::<Vec<_>>()
1094 .join(",")
1095 })
1096 .unwrap_or_default()
1097 };
1098 let hm = || {
1099 format!(
1100 "{:02}:{:02}",
1101 trigger.get("hour").and_then(|v| v.as_u64()).unwrap_or(0),
1102 trigger.get("minute").and_then(|v| v.as_u64()).unwrap_or(0)
1103 )
1104 };
1105 match trigger.get("type").and_then(|t| t.as_str()) {
1106 Some("interval") => format!(
1107 "every {}s",
1108 trigger
1109 .get("every_seconds")
1110 .and_then(|v| v.as_u64())
1111 .unwrap_or(0)
1112 ),
1113 Some("daily") => format!(
1114 "daily {}:{:02}",
1115 hm(),
1116 trigger.get("second").and_then(|v| v.as_u64()).unwrap_or(0)
1117 ),
1118 Some("weekly") => format!("weekly {} {}", joined("weekdays"), hm()),
1119 Some("monthly") => format!("monthly {} {}", joined("days"), hm()),
1120 Some("cron") => format!(
1121 "cron '{}'",
1122 trigger.get("expr").and_then(|v| v.as_str()).unwrap_or("?")
1123 ),
1124 _ => trigger.to_string(),
1125 }
1126}
1127
1128fn fmt_ts(value: Option<&serde_json::Value>) -> String {
1131 let Some(s) = value.and_then(|v| v.as_str()) else {
1132 return "-".to_string();
1133 };
1134 if s.len() >= 19 && s.is_char_boundary(19) && s.as_bytes().get(10) == Some(&b'T') {
1136 format!("{} {}", &s[..10], &s[11..19])
1137 } else {
1138 s.to_string()
1139 }
1140}
1141
1142const MCP_MUTATE_TIMEOUT: Duration = Duration::from_secs(60);
1152
1153pub async fn mcp_status(conn: ConnArgs, json: bool) -> anyhow::Result<()> {
1156 let base = conn.api_base();
1157 let url = format!("{base}/mcp/servers");
1158 let resp = reqwest::Client::new()
1159 .get(&url)
1160 .timeout(REQUEST_TIMEOUT)
1161 .send()
1162 .await
1163 .map_err(|e| unreachable(&base, e))?;
1164 if !resp.status().is_success() {
1165 anyhow::bail!("GET {url} -> HTTP {}", resp.status());
1166 }
1167 let v: serde_json::Value = resp.json().await?;
1168 if json {
1169 println!("{}", serde_json::to_string_pretty(&v)?);
1170 return Ok(());
1171 }
1172
1173 let servers = v.get("servers").and_then(|s| s.as_array());
1174 let servers = match servers {
1175 Some(s) if !s.is_empty() => s,
1176 _ => {
1177 println!("(no MCP servers configured)");
1178 return Ok(());
1179 }
1180 };
1181
1182 println!(
1184 "{:<24} {:<8} {:<14} {:>5} NAME",
1185 "ID", "ENABLED", "STATUS", "TOOLS"
1186 );
1187 for s in servers {
1188 let id = s.get("id").and_then(|x| x.as_str()).unwrap_or("?");
1189 let enabled = s.get("enabled").and_then(|b| b.as_bool()).unwrap_or(false);
1190 let status = s.get("status").and_then(|x| x.as_str()).unwrap_or("?");
1191 let tools = s.get("tool_count").and_then(|x| x.as_u64()).unwrap_or(0);
1192 let name = s.get("name").and_then(|x| x.as_str()).unwrap_or("");
1193 println!(
1194 "{:<24} {:<8} {:<14} {:>5} {}",
1195 truncate(id, 24),
1196 if enabled { "yes" } else { "no" },
1197 truncate(status, 14),
1198 tools,
1199 truncate(name, 40)
1200 );
1201 }
1202 for s in servers {
1203 if let Some(err) = s.get("last_error").and_then(|e| e.as_str()) {
1204 if !err.trim().is_empty() {
1205 let id = s.get("id").and_then(|x| x.as_str()).unwrap_or("?");
1206 println!("{} {id}: {}", "!".yellow(), truncate(err, 100));
1207 }
1208 }
1209 }
1210 let connected = servers
1211 .iter()
1212 .filter(|s| s.get("status").and_then(|x| x.as_str()) == Some("connected"))
1213 .count();
1214 println!("\n{connected} connected of {}.", servers.len());
1215 Ok(())
1216}
1217
1218pub async fn mcp_connect(conn: ConnArgs, server_id: &str) -> anyhow::Result<()> {
1221 guard_id_segment("MCP server id", server_id)?;
1222 let base = conn.api_base();
1223 let url = format!("{base}/mcp/servers/{server_id}/connect");
1224 let resp = reqwest::Client::new()
1225 .post(&url)
1226 .timeout(MCP_MUTATE_TIMEOUT)
1227 .send()
1228 .await
1229 .map_err(|e| unreachable(&base, e))?;
1230 let status = resp.status();
1231 let body: serde_json::Value = resp.json().await.unwrap_or(serde_json::Value::Null);
1232 if status.is_success() {
1233 println!("{} server '{server_id}' connected", "✓".green());
1234 Ok(())
1235 } else if status.as_u16() == 404 {
1236 anyhow::bail!("MCP server '{server_id}' not found (check `bamboo mcp status`)");
1237 } else {
1238 anyhow::bail!(
1239 "connect failed: HTTP {status} {}",
1240 server_error_message(&body)
1241 );
1242 }
1243}
1244
1245pub async fn mcp_disconnect(conn: ConnArgs, server_id: &str) -> anyhow::Result<()> {
1248 guard_id_segment("MCP server id", server_id)?;
1249 let base = conn.api_base();
1250 let url = format!("{base}/mcp/servers/{server_id}/disconnect");
1251 let resp = reqwest::Client::new()
1252 .post(&url)
1253 .timeout(MCP_MUTATE_TIMEOUT)
1254 .send()
1255 .await
1256 .map_err(|e| unreachable(&base, e))?;
1257 let status = resp.status();
1258 let body: serde_json::Value = resp.json().await.unwrap_or(serde_json::Value::Null);
1259 if status.is_success() {
1260 println!("{} server '{server_id}' disconnected", "✓".green());
1261 Ok(())
1262 } else if status.as_u16() == 404 {
1263 anyhow::bail!("MCP server '{server_id}' not found (check `bamboo mcp status`)");
1264 } else {
1265 anyhow::bail!(
1266 "disconnect failed: HTTP {status} {}",
1267 server_error_message(&body)
1268 );
1269 }
1270}
1271
1272pub async fn mcp_refresh(conn: ConnArgs, server_id: Option<&str>) -> anyhow::Result<()> {
1275 let base = conn.api_base();
1276 let client = reqwest::Client::new();
1277
1278 let targets: Vec<String> = match server_id {
1279 Some(id) => {
1280 guard_id_segment("MCP server id", id)?;
1281 vec![id.to_string()]
1282 }
1283 None => {
1284 let url = format!("{base}/mcp/servers");
1286 let resp = client
1287 .get(&url)
1288 .timeout(REQUEST_TIMEOUT)
1289 .send()
1290 .await
1291 .map_err(|e| unreachable(&base, e))?;
1292 if !resp.status().is_success() {
1293 anyhow::bail!("GET {url} -> HTTP {}", resp.status());
1294 }
1295 let v: serde_json::Value = resp.json().await?;
1296 let ids: Vec<String> = v
1297 .get("servers")
1298 .and_then(|s| s.as_array())
1299 .map(|servers| {
1300 servers
1301 .iter()
1302 .filter(|s| s.get("enabled").and_then(|b| b.as_bool()).unwrap_or(false))
1303 .filter_map(|s| s.get("id").and_then(|x| x.as_str()))
1304 .map(String::from)
1305 .collect()
1306 })
1307 .unwrap_or_default();
1308 if ids.is_empty() {
1309 println!("(no enabled MCP servers to refresh)");
1310 return Ok(());
1311 }
1312 ids
1313 }
1314 };
1315
1316 let mut failures = 0usize;
1317 for id in &targets {
1318 let url = format!("{base}/mcp/servers/{id}/refresh");
1319 let result = client.post(&url).timeout(MCP_MUTATE_TIMEOUT).send().await;
1320 match result {
1321 Ok(resp) if resp.status().is_success() => {
1322 let body: serde_json::Value = resp.json().await.unwrap_or(serde_json::Value::Null);
1323 let tools = body.get("tool_count").and_then(|t| t.as_u64()).unwrap_or(0);
1324 println!("{} {id}: {tools} tool(s)", "✓".green());
1325 }
1326 Ok(resp) => {
1327 let status = resp.status();
1328 let body: serde_json::Value = resp.json().await.unwrap_or(serde_json::Value::Null);
1329 println!(
1330 "{} {id}: HTTP {status} {}",
1331 "✗".red(),
1332 server_error_message(&body)
1333 );
1334 failures += 1;
1335 }
1336 Err(e) => {
1337 println!("{} {id}: {e}", "✗".red());
1338 failures += 1;
1339 }
1340 }
1341 }
1342 if failures > 0 {
1343 anyhow::bail!("{failures} of {} refresh(es) failed", targets.len());
1344 }
1345 Ok(())
1346}
1347
1348pub async fn mcp_tools(conn: ConnArgs, server_id: Option<&str>, json: bool) -> anyhow::Result<()> {
1352 let base = conn.api_base();
1353 let url = match server_id {
1354 Some(id) => {
1355 guard_id_segment("MCP server id", id)?;
1356 format!("{base}/mcp/servers/{id}/tools")
1357 }
1358 None => format!("{base}/mcp/tools"),
1359 };
1360 let resp = reqwest::Client::new()
1361 .get(&url)
1362 .timeout(REQUEST_TIMEOUT)
1363 .send()
1364 .await
1365 .map_err(|e| unreachable(&base, e))?;
1366 if resp.status().as_u16() == 404 {
1367 anyhow::bail!(
1368 "MCP server '{}' not found (check `bamboo mcp status`)",
1369 server_id.unwrap_or("?")
1370 );
1371 }
1372 if !resp.status().is_success() {
1373 anyhow::bail!("GET {url} -> HTTP {}", resp.status());
1374 }
1375 let v: serde_json::Value = resp.json().await?;
1376 if json {
1377 println!("{}", serde_json::to_string_pretty(&v)?);
1378 return Ok(());
1379 }
1380
1381 let tools = v.get("tools").and_then(|t| t.as_array());
1382 let tools = match tools {
1383 Some(t) if !t.is_empty() => t,
1384 _ => {
1385 println!("(no tools)");
1386 return Ok(());
1387 }
1388 };
1389 println!("{:<32} {:<20} DESCRIPTION", "ALIAS", "SERVER");
1390 for t in tools {
1391 let alias = t.get("alias").and_then(|x| x.as_str()).unwrap_or("?");
1392 let server = t.get("server_id").and_then(|x| x.as_str()).unwrap_or("");
1393 let desc = t.get("description").and_then(|x| x.as_str()).unwrap_or("");
1394 println!(
1395 "{:<32} {:<20} {}",
1396 truncate(alias, 32),
1397 truncate(server, 20),
1398 truncate(desc, 70)
1399 );
1400 }
1401 println!("\n{} tool(s).", tools.len());
1402 Ok(())
1403}
1404
1405pub async fn mcp_add(conn: ConnArgs, payload_source: &str) -> anyhow::Result<()> {
1410 let payload = read_json_payload(payload_source)?;
1412
1413 let base = conn.api_base();
1414 let url = format!("{base}/mcp/servers");
1415 let resp = reqwest::Client::new()
1416 .post(&url)
1417 .timeout(MCP_MUTATE_TIMEOUT)
1418 .json(&payload)
1419 .send()
1420 .await
1421 .map_err(|e| unreachable(&base, e))?;
1422 let status = resp.status();
1423 let body: serde_json::Value = resp.json().await.unwrap_or(serde_json::Value::Null);
1424 if status.is_success() {
1425 let id = body
1426 .get("server_id")
1427 .and_then(|s| s.as_str())
1428 .unwrap_or("?");
1429 println!("{} server '{id}' saved", "✓".green());
1430 Ok(())
1431 } else {
1432 anyhow::bail!("add failed: HTTP {status} {}", server_error_message(&body));
1433 }
1434}
1435
1436pub async fn mcp_remove(conn: ConnArgs, server_id: &str, yes: bool) -> anyhow::Result<()> {
1441 guard_id_segment("MCP server id", server_id)?;
1442 if !yes
1443 && !confirm(&format!(
1444 "Remove MCP server '{server_id}'? This stops it and deletes its stored config."
1445 ))?
1446 {
1447 println!("aborted (nothing removed).");
1448 return Ok(());
1449 }
1450 let base = conn.api_base();
1451 let url = format!("{base}/mcp/servers/{server_id}");
1452 let resp = reqwest::Client::new()
1453 .delete(&url)
1454 .timeout(MCP_MUTATE_TIMEOUT)
1455 .send()
1456 .await
1457 .map_err(|e| unreachable(&base, e))?;
1458 let status = resp.status();
1459 let body: serde_json::Value = resp.json().await.unwrap_or(serde_json::Value::Null);
1460 if status.is_success() {
1461 println!("{} server '{server_id}' removed", "✓".green());
1462 Ok(())
1463 } else {
1464 anyhow::bail!(
1465 "remove failed: HTTP {status} {}",
1466 server_error_message(&body)
1467 );
1468 }
1469}
1470
1471fn count_running(sessions: &[serde_json::Value]) -> usize {
1473 sessions
1474 .iter()
1475 .filter(|x| {
1476 x.get("is_running")
1477 .and_then(|b| b.as_bool())
1478 .unwrap_or(false)
1479 })
1480 .count()
1481}
1482
1483pub(crate) fn truncate(s: &str, max: usize) -> String {
1485 if s.chars().count() <= max {
1486 s.to_string()
1487 } else {
1488 let head: String = s.chars().take(max.saturating_sub(1)).collect();
1489 format!("{head}…")
1490 }
1491}
1492
1493#[cfg(test)]
1494mod tests {
1495 use super::{
1496 build_create_payload, fmt_ts, format_pending_question, format_session_detail,
1497 guard_id_segment, history_summary, parse_daily_time, server_error_message, trigger_summary,
1498 ScheduleCreateArgs,
1499 };
1500
1501 #[test]
1502 fn guard_id_segment_rejects_path_hazards() {
1503 for bad in [
1504 "", ".", "..", "a/b", "a\\b", "a?b", "a#b", "a%b", "a b", "a\tb",
1505 ] {
1506 assert!(
1507 guard_id_segment("session id", bad).is_err(),
1508 "{bad:?} must be rejected"
1509 );
1510 }
1511 assert!(guard_id_segment("session id", "0195fd1e-abc4-7def-8123-456789abcdef").is_ok());
1512 let err = guard_id_segment("schedule id", "a/b").unwrap_err();
1514 assert!(err.to_string().contains("invalid schedule id"));
1515 }
1516
1517 #[test]
1518 fn history_summary_reports_true_total_and_truncation() {
1519 assert_eq!(
1521 history_summary("s1", 3, 3, false),
1522 "3 message(s) in session s1."
1523 );
1524 let line = history_summary("s1", 2000, 5000, true);
1526 assert!(line.starts_with("5000 message(s) in session s1"));
1527 assert!(line.contains("newest 2000"));
1528 assert!(line.contains("history cap"));
1529 }
1530
1531 #[test]
1532 fn format_pending_question_pretty_prints_question_and_options() {
1533 colored::control::set_override(false);
1534 let v = serde_json::json!({
1535 "has_pending_question": true,
1536 "question": "Proceed with the deploy?",
1537 "options": ["Yes", "No"],
1538 "allow_custom": true,
1539 "tool_call_id": "tc-1",
1540 "tool_name": "conclusion_with_options",
1541 "source": "tool",
1542 });
1543 let text = format_pending_question("sess-1", &v).expect("pending question");
1544 assert!(text.contains("question: Proceed with the deploy?"));
1545 assert!(text.contains("1. Yes"));
1546 assert!(text.contains("2. No"));
1547 assert!(text.contains("custom free-text answers are allowed"));
1548 assert!(text.contains("tool: conclusion_with_options"));
1549 assert!(text.contains("bamboo respond sess-1"));
1550 assert!(text.contains("resumes the run server-side"));
1551 }
1552
1553 #[test]
1554 fn format_pending_question_none_when_no_question() {
1555 let v = serde_json::json!({ "has_pending_question": false });
1556 assert!(format_pending_question("sess-1", &v).is_none());
1557 assert!(format_pending_question("sess-1", &serde_json::json!({})).is_none());
1559 }
1560
1561 #[test]
1562 fn format_session_detail_shows_core_and_optional_fields() {
1563 colored::control::set_override(false);
1564 let s = serde_json::json!({
1565 "id": "sess-9",
1566 "title": "Fix the bug",
1567 "kind": "root",
1568 "model": "claude-sonnet-5",
1569 "provider": "anthropic",
1570 "is_running": true,
1571 "last_run_status": "failed",
1572 "last_run_error": "boom",
1573 "has_pending_question": true,
1574 "message_count": 12,
1575 "pinned": true,
1576 "running_child_count": 2,
1577 "created_at": "2026-07-10T00:00:00Z",
1578 "last_activity_at": "2026-07-10T01:00:00Z",
1579 "placement": { "kind": "local", "host": "mac.local" },
1580 });
1581 let text = format_session_detail(&s);
1582 assert!(text.contains("sess-9"));
1583 assert!(text.contains("Fix the bug"));
1584 assert!(text.contains("anthropic:claude-sonnet-5"));
1585 assert!(text.contains("failed (boom)"));
1586 assert!(text.contains("2 running"));
1587 assert!(text.contains("local @ mac.local"));
1588 assert!(text.contains("bamboo respond <id> --pending"));
1589
1590 let minimal = serde_json::json!({
1592 "id": "sess-min",
1593 "title": "",
1594 "kind": "root",
1595 "model": "m",
1596 "is_running": false,
1597 "message_count": 0,
1598 });
1599 let text = format_session_detail(&minimal);
1600 assert!(!text.contains("last run"));
1601 assert!(!text.contains("parent"));
1602 assert!(!text.contains("children"));
1603 assert!(!text.contains("placement"));
1604 }
1605 #[test]
1606 fn parse_daily_time_accepts_hm_and_hms() {
1607 assert_eq!(parse_daily_time("09:30").unwrap(), (9, 30, 0));
1608 assert_eq!(parse_daily_time("23:59:59").unwrap(), (23, 59, 59));
1609 }
1610
1611 #[test]
1612 fn parse_daily_time_rejects_malformed_and_out_of_range() {
1613 for bad in ["", "9", "24:00", "09:60", "09:30:60", "a:b", "09:30:15:00"] {
1614 assert!(parse_daily_time(bad).is_err(), "{bad:?}");
1615 }
1616 }
1617
1618 #[test]
1619 fn build_create_payload_maps_flags_to_request_shape() {
1620 let payload = build_create_payload(&ScheduleCreateArgs {
1621 name: Some("nightly".to_string()),
1622 cron: Some("0 0 2 * * *".to_string()),
1623 prompt: Some("run the suite".to_string()),
1624 model: Some("anthropic:claude-sonnet-4".to_string()),
1625 workspace: Some("/tmp/repo".to_string()),
1626 timezone: Some("Asia/Shanghai".to_string()),
1627 ..Default::default()
1628 })
1629 .unwrap();
1630
1631 assert_eq!(payload["name"], "nightly");
1632 assert_eq!(payload["enabled"], true);
1633 assert_eq!(payload["trigger"]["type"], "cron");
1634 assert_eq!(payload["trigger"]["expr"], "0 0 2 * * *");
1635 assert_eq!(payload["timezone"], "Asia/Shanghai");
1636 assert_eq!(payload["run_config"]["task_message"], "run the suite");
1637 assert_eq!(payload["run_config"]["auto_execute"], true);
1638 assert_eq!(payload["run_config"]["model"], "anthropic:claude-sonnet-4");
1639 assert_eq!(payload["run_config"]["workspace_path"], "/tmp/repo");
1640 }
1641
1642 #[test]
1643 fn build_create_payload_daily_and_disabled() {
1644 let payload = build_create_payload(&ScheduleCreateArgs {
1645 name: Some("standup".to_string()),
1646 daily: Some("09:30".to_string()),
1647 prompt: Some("summarize".to_string()),
1648 disabled: true,
1649 ..Default::default()
1650 })
1651 .unwrap();
1652
1653 assert_eq!(payload["enabled"], false);
1654 assert_eq!(payload["trigger"]["type"], "daily");
1655 assert_eq!(payload["trigger"]["hour"], 9);
1656 assert_eq!(payload["trigger"]["minute"], 30);
1657 assert_eq!(payload["trigger"]["second"], 0);
1658 assert!(payload.get("timezone").is_none());
1660 assert!(payload["run_config"].get("model").is_none());
1661 }
1662
1663 #[test]
1664 fn build_create_payload_interval_trigger() {
1665 let payload = build_create_payload(&ScheduleCreateArgs {
1666 name: Some("tick".to_string()),
1667 every: Some(3600),
1668 prompt: Some("check".to_string()),
1669 ..Default::default()
1670 })
1671 .unwrap();
1672 assert_eq!(payload["trigger"]["type"], "interval");
1673 assert_eq!(payload["trigger"]["every_seconds"], 3600);
1674 }
1675
1676 #[test]
1677 fn build_create_payload_requires_name_prompt_and_trigger() {
1678 assert!(build_create_payload(&ScheduleCreateArgs::default()).is_err());
1679 assert!(build_create_payload(&ScheduleCreateArgs {
1680 name: Some("x".to_string()),
1681 prompt: Some("y".to_string()),
1682 ..Default::default()
1683 })
1684 .is_err());
1685 }
1686
1687 #[test]
1688 fn trigger_summary_renders_each_kind() {
1689 let case = |json: serde_json::Value| trigger_summary(&json);
1690 assert_eq!(
1691 case(serde_json::json!({"type":"interval","every_seconds":60})),
1692 "every 60s"
1693 );
1694 assert_eq!(
1695 case(serde_json::json!({"type":"daily","hour":9,"minute":30,"second":0})),
1696 "daily 09:30:00"
1697 );
1698 assert_eq!(
1699 case(serde_json::json!({"type":"weekly","weekdays":["mon","fri"],"hour":9,"minute":0})),
1700 "weekly mon,fri 09:00"
1701 );
1702 assert_eq!(
1703 case(serde_json::json!({"type":"monthly","days":[1,15],"hour":8,"minute":5})),
1704 "monthly 1,15 08:05"
1705 );
1706 assert_eq!(
1707 case(serde_json::json!({"type":"cron","expr":"0 0 2 * * *"})),
1708 "cron '0 0 2 * * *'"
1709 );
1710 }
1711
1712 #[test]
1713 fn fmt_ts_shortens_rfc3339_and_defaults_to_dash() {
1714 let value = serde_json::json!("2026-07-10T12:34:56.789012Z");
1715 assert_eq!(fmt_ts(Some(&value)), "2026-07-10 12:34:56");
1716 assert_eq!(fmt_ts(None), "-");
1717 let null = serde_json::Value::Null;
1718 assert_eq!(fmt_ts(Some(&null)), "-");
1719 }
1720
1721 #[test]
1722 fn server_error_message_extracts_nested_and_legacy_error_fields() {
1723 assert_eq!(
1724 server_error_message(&serde_json::json!({
1725 "error": {"message": "name is required", "type": "api_error"}
1726 })),
1727 "(name is required)"
1728 );
1729 assert_eq!(
1730 server_error_message(&serde_json::json!({"error":"name is required"})),
1731 "(name is required)"
1732 );
1733 assert_eq!(server_error_message(&serde_json::Value::Null), "");
1734 }
1735}