1use std::net::SocketAddr;
11use std::path::PathBuf;
12use std::time::Duration;
13
14use chrono::{Duration as ChronoDuration, Utc};
15
16use bamboo_llm::Config;
17use bamboo_subagent::discovery::Fabric;
18use bamboo_subagent::executor::{ChildExecutor, EchoExecutor};
19use bamboo_subagent::fleet::spawn_worker;
20use bamboo_subagent::proto::{AgentRecord, ChildFrame, ParentFrame, RunSpec, TerminalStatus};
21use bamboo_subagent::provision::{
22 ChildIdentity, ExecutorSpec, ModelRefSpec, ProvisionSpec, ScopedCredential,
23};
24use bamboo_subagent::transport::{ChildClient, WsServer};
25
26use crate::subagent_worker::BambooRuntimeExecutor;
27
28pub fn default_fabric_dir() -> PathBuf {
34 bamboo_config::paths::subagents_dir()
35}
36
37pub struct ActorRunArgs {
38 pub prompt: String,
39 pub model: Option<String>,
40 pub role: String,
41 pub workspace: Option<PathBuf>,
42 pub data_dir: Option<PathBuf>,
43 pub echo: bool,
44 pub raw: bool,
46}
47
48pub struct ActorServeArgs {
49 pub role: String,
50 pub id: Option<String>,
52 pub model: Option<String>,
53 pub workspace: Option<PathBuf>,
54 pub data_dir: Option<PathBuf>,
55 pub echo: bool,
56 pub bind: Option<SocketAddr>,
60 pub tls: bool,
62 pub cert_file: Option<PathBuf>,
63 pub key_file: Option<PathBuf>,
64 pub token: Option<String>,
67}
68
69pub struct ActorCallArgs {
70 pub agent: String,
72 pub prompt: String,
73 pub raw: bool,
74}
75
76pub async fn run(args: ActorRunArgs) -> Result<(), String> {
81 let child_id = format!("cli-{}", uuid::Uuid::new_v4());
82 let spec = prepare_spec(
83 &child_id,
84 &args.role,
85 &args.model,
86 &args.workspace,
87 &args.data_dir,
88 args.echo,
89 )?;
90
91 let worker_bin =
92 std::env::current_exe().map_err(|e| format!("cannot locate own executable: {e}"))?;
93 eprintln!(
94 "▶ spawning actor {child_id} (model: {}, executor: {})",
95 describe_model(&spec),
96 if args.echo { "echo" } else { "bamboo_runtime" },
97 );
98
99 let spawned = spawn_worker(
100 &worker_bin,
101 &["subagent-worker".to_string()],
102 &spec,
103 Duration::from_secs(30),
104 )
105 .await
106 .map_err(|e| format!("spawn/register failed: {e}"))?;
107 eprintln!(
108 "✔ actor registered (pid {}, endpoint {})",
109 spawned.record.pid, spawned.record.endpoint
110 );
111
112 let exit = connect_and_stream(&spawned.record.endpoint, &args.prompt, args.raw).await;
113 spawned.kill().await;
114 exit
115}
116
117pub async fn serve(args: ActorServeArgs) -> Result<(), String> {
122 let agent_id = args
123 .id
124 .clone()
125 .unwrap_or_else(|| format!("{}-{}", args.role, &uuid::Uuid::new_v4().to_string()[..8]));
126 let spec = prepare_spec(
127 &agent_id,
128 &args.role,
129 &args.model,
130 &args.workspace,
131 &args.data_dir,
132 args.echo,
133 )?;
134
135 let executor: std::sync::Arc<dyn ChildExecutor> = if args.echo {
136 std::sync::Arc::new(EchoExecutor)
137 } else {
138 std::sync::Arc::new(BambooRuntimeExecutor::build(&spec).await?)
139 };
140
141 let server = if args.tls {
144 let (cert, key) = match (&args.cert_file, &args.key_file) {
145 (Some(c), Some(k)) => (c, k),
146 _ => return Err("--tls requires both --cert-file and --key-file".to_string()),
147 };
148 let bind_addr = args
149 .bind
150 .unwrap_or_else(|| (std::net::Ipv4Addr::UNSPECIFIED, 8443).into());
151 WsServer::bind_tls(bind_addr, cert, key, args.token.clone())
152 .await
153 .map_err(|e| format!("bind_tls: {e}"))?
154 } else if let Some(bind_addr) = args.bind {
155 if args.token.is_some() && !bind_addr.ip().is_loopback() {
160 return Err(format!(
161 "refusing --token on a non-loopback plaintext bind ({bind_addr}): the token would \
162 be sent in cleartext. Use --tls (with --cert-file/--key-file) for a public bind."
163 ));
164 }
165 WsServer::bind_with_token(bind_addr, args.token.clone())
166 .await
167 .map_err(|e| format!("bind: {e}"))?
168 } else {
169 WsServer::bind_loopback()
171 .await
172 .map_err(|e| format!("bind: {e}"))?
173 };
174 let endpoint = server.ws_endpoint();
175
176 let fab = std::sync::Arc::new(Fabric::at(&spec.fabric_dir));
177 let _ = fab.gc().await; let record = AgentRecord {
179 agent_id: agent_id.clone(),
180 role: args.role.clone(),
181 labels: Vec::new(),
182 endpoint: endpoint.clone(),
183 pid: std::process::id(),
184 version: env!("CARGO_PKG_VERSION").to_string(),
185 started_at: Utc::now(),
186 lease_expires_at: Utc::now() + ChronoDuration::seconds(60),
187 };
188 fab.publish(&record)
189 .await
190 .map_err(|e| format!("announce: {e}"))?;
191
192 let renew_fab = fab.clone();
194 let mut renew_record = record.clone();
195 let renew = tokio::spawn(async move {
196 let mut tick = tokio::time::interval(Duration::from_secs(20));
197 tick.tick().await;
198 loop {
199 tick.tick().await;
200 renew_record.lease_expires_at = Utc::now() + ChronoDuration::seconds(60);
201 if renew_fab.publish(&renew_record).await.is_err() {
202 break;
203 }
204 }
205 });
206
207 eprintln!(
208 "✔ service agent '{agent_id}' (role: {}) announced at {endpoint}",
209 args.role
210 );
211 eprintln!(" serving until Ctrl-C — call it with: bamboo actor call {agent_id} \"<task>\"");
212
213 let result = tokio::select! {
215 r = server.serve(executor) => r.map_err(|e| format!("serve: {e}")),
216 _ = tokio::signal::ctrl_c() => Ok(()),
217 };
218 renew.abort();
219 let _ = fab.withdraw(&agent_id).await;
220 eprintln!("⏹ service agent '{agent_id}' withdrawn");
221 result
222}
223
224pub async fn list() -> Result<(), String> {
229 let fab = Fabric::at(default_fabric_dir());
230 let _ = fab.gc().await;
231 let records = fab.discover().await.map_err(|e| format!("discover: {e}"))?;
232 if records.is_empty() {
233 println!(
234 "no live actors (fabric: {})",
235 default_fabric_dir().display()
236 );
237 return Ok(());
238 }
239 println!("{:<28} {:<12} {:<8} ENDPOINT", "AGENT", "ROLE", "PID");
240 for r in records {
241 println!(
242 "{:<28} {:<12} {:<8} {}",
243 r.agent_id, r.role, r.pid, r.endpoint
244 );
245 }
246 Ok(())
247}
248
249pub async fn call(args: ActorCallArgs) -> Result<(), String> {
254 let fab = Fabric::at(default_fabric_dir());
255 let record = match fab
256 .resolve(&args.agent)
257 .await
258 .map_err(|e| format!("resolve: {e}"))?
259 {
260 Some(r) => r,
261 None => {
262 fab.discover()
264 .await
265 .map_err(|e| format!("discover: {e}"))?
266 .into_iter()
267 .find(|r| r.role == args.agent)
268 .ok_or_else(|| {
269 format!(
270 "no live actor with id or role '{}'; see `bamboo actor list`",
271 args.agent
272 )
273 })?
274 }
275 };
276 eprintln!(
277 "▶ calling {} (role: {}, endpoint {})",
278 record.agent_id, record.role, record.endpoint
279 );
280 connect_and_stream(&record.endpoint, &args.prompt, args.raw).await
281}
282
283async fn connect_and_stream(endpoint: &str, prompt: &str, raw: bool) -> Result<(), String> {
290 let mut client = ChildClient::connect(endpoint)
291 .await
292 .map_err(|e| format!("connect failed: {e}"))?;
293 client
294 .send(ParentFrame::Run(RunSpec {
295 assignment: prompt.to_string(),
296 reasoning_effort: None,
297 messages: Vec::new(),
298 }))
299 .await
300 .map_err(|e| format!("dispatch failed: {e}"))?;
301
302 let (cancel_tx, mut cancel_rx) = tokio::sync::mpsc::channel::<()>(1);
303 tokio::spawn(async move {
304 if tokio::signal::ctrl_c().await.is_ok() {
305 let _ = cancel_tx.send(()).await;
306 }
307 });
308
309 let mut exit: Result<(), String> = Ok(());
310 let mut streamed_tokens = false;
311 loop {
312 tokio::select! {
313 _ = cancel_rx.recv() => {
314 eprintln!("\n⏹ cancelling…");
315 let _ = client.send(ParentFrame::Cancel).await;
316 }
317 frame = client.next_frame() => {
318 match frame {
319 Ok(Some(ChildFrame::Event { event })) => {
320 if event["type"] == "token" {
321 streamed_tokens = true;
322 }
323 print_event(&event, raw);
324 }
325 Ok(Some(ChildFrame::ApprovalRequest { .. })) => {
326 }
329 Ok(Some(ChildFrame::Terminal { status, result, error, .. })) => {
330 println!();
331 match status {
332 TerminalStatus::Completed => {
333 eprintln!("✔ completed");
334 if !streamed_tokens {
335 if let Some(r) = result {
336 println!("{r}");
337 }
338 }
339 }
340 TerminalStatus::Cancelled => eprintln!("⏹ cancelled"),
341 TerminalStatus::Suspended => eprintln!("⏸ suspended (waiting on sub-agents)"),
342 TerminalStatus::Error => {
343 exit = Err(error.unwrap_or_else(|| "actor errored".into()));
344 }
345 }
346 break;
347 }
348 Ok(None) => {
349 exit = Err("connection closed before terminal".into());
350 break;
351 }
352 Err(e) => {
353 exit = Err(format!("transport error: {e}"));
354 break;
355 }
356 }
357 }
358 }
359 }
360
361 let _ = client.close().await;
362 exit
363}
364
365fn prepare_spec(
367 child_id: &str,
368 role: &str,
369 model_arg: &Option<String>,
370 workspace: &Option<PathBuf>,
371 data_dir: &Option<PathBuf>,
372 echo: bool,
373) -> Result<ProvisionSpec, String> {
374 let data_dir = data_dir
375 .clone()
376 .unwrap_or_else(bamboo_config::paths::resolve_bamboo_dir);
377 let config = Config::from_data_dir(Some(data_dir.clone()));
379 let credentials =
380 bamboo_engine::external_agents::runtime::extract_provider_credentials(&config);
381
382 let model = resolve_model(model_arg, &config)?;
383 if !echo && model.is_none() {
384 return Err(
385 "no model resolved: pass --model provider:model or configure defaults.sub_agent/chat"
386 .to_string(),
387 );
388 }
389
390 let mut spec = ProvisionSpec::new(
391 ChildIdentity {
392 child_id: child_id.to_string(),
393 parent_id: None,
394 project_key: None,
395 role: role.to_string(),
396 depth: 0,
397 },
398 if echo {
399 ExecutorSpec::Echo
400 } else {
401 ExecutorSpec::BambooRuntime
402 },
403 default_fabric_dir().to_string_lossy().into_owned(),
404 );
405 spec.workspace = workspace
406 .clone()
407 .or_else(|| std::env::current_dir().ok())
408 .map(|w| w.to_string_lossy().into_owned());
409 spec.model = model.clone();
410 if let Some(m) = &model {
411 if let Some(cred) = pick_credential(&credentials, &m.provider) {
412 spec.secrets.provider_credentials.push(cred);
413 } else if !echo {
414 return Err(format!(
415 "no credential found for provider '{}' in {}",
416 m.provider,
417 data_dir.display()
418 ));
419 }
420 }
421 Ok(spec)
422}
423
424fn describe_model(spec: &ProvisionSpec) -> String {
425 spec.model
426 .as_ref()
427 .map(|m| format!("{}:{}", m.provider, m.model))
428 .unwrap_or_else(|| "-".into())
429}
430
431fn print_event(event: &serde_json::Value, raw: bool) {
432 use std::io::Write;
433 if raw {
434 println!("{event}");
435 return;
436 }
437 match event["type"].as_str().unwrap_or("") {
438 "token" => {
439 print!("{}", event["content"].as_str().unwrap_or(""));
440 let _ = std::io::stdout().flush();
441 }
442 "reasoning_token" => { }
443 "tool_start" => {
444 eprintln!("\n⚙ {}", event["tool_name"].as_str().unwrap_or("tool"));
445 }
446 "tool_complete" => eprintln!("✔ tool done"),
447 "tool_error" => eprintln!("✘ tool error: {}", event["error"].as_str().unwrap_or("")),
448 "error" => eprintln!("✘ {}", event["message"].as_str().unwrap_or("")),
449 _ => {}
450 }
451}
452
453fn resolve_model(
458 explicit: &Option<String>,
459 config: &Config,
460) -> Result<Option<ModelRefSpec>, String> {
461 if let Some(raw) = explicit {
462 if let Some(parsed) =
463 crate::model_spec::parse_model_spec(raw).map_err(|e| format!("--model {e}"))?
464 {
465 let provider = parsed.provider.unwrap_or_else(|| config.provider.clone());
466 return Ok(Some(ModelRefSpec {
467 provider,
468 model: parsed.model,
469 }));
470 }
471 }
472 if let Some(defaults) = &config.defaults {
473 let pick = defaults.sub_agent.as_ref().or(Some(&defaults.chat));
474 if let Some(r) = pick {
475 return Ok(Some(ModelRefSpec {
476 provider: r.provider.clone(),
477 model: r.model.clone(),
478 }));
479 }
480 }
481 Ok(None)
482}
483
484fn pick_credential(creds: &[ScopedCredential], provider: &str) -> Option<ScopedCredential> {
485 creds.iter().find(|c| c.provider == provider).cloned()
486}
487
488#[cfg(test)]
489mod tests {
490 use super::*;
491
492 fn some(s: &str) -> Option<String> {
493 Some(s.to_string())
494 }
495
496 #[test]
499 fn resolve_model_colon_form() {
500 let config = Config {
501 provider: "anthropic".into(),
502 ..Config::default()
503 };
504 let m = resolve_model(&some("openai:gpt-4o"), &config)
505 .unwrap()
506 .unwrap();
507 assert_eq!(m.provider, "openai");
508 assert_eq!(m.model, "gpt-4o");
509 }
510
511 #[test]
514 fn resolve_model_bare_uses_config_default_provider() {
515 let config = Config {
516 provider: "openai".into(),
517 ..Config::default()
518 };
519 let m = resolve_model(&some("gpt-4o"), &config).unwrap().unwrap();
520 assert_eq!(m.provider, "openai");
521 assert_eq!(m.model, "gpt-4o");
522 }
523
524 fn defaults_config(
525 chat: (&str, &str),
526 sub_agent: Option<(&str, &str)>,
527 ) -> bamboo_config::DefaultsConfig {
528 bamboo_config::DefaultsConfig {
529 chat: bamboo_domain::ProviderModelRef::new(chat.0, chat.1),
530 fast: None,
531 task_summary: None,
532 vision: None,
533 memory_background: None,
534 planning: None,
535 search: None,
536 code_review: None,
537 sub_agent: sub_agent.map(|(p, m)| bamboo_domain::ProviderModelRef::new(p, m)),
538 subagent_models: std::collections::HashMap::new(),
539 }
540 }
541
542 #[test]
544 fn resolve_model_falls_back_to_defaults_sub_agent() {
545 let config = Config {
546 defaults: Some(defaults_config(
547 ("anthropic", "claude-x"),
548 Some(("openai", "gpt-sub")),
549 )),
550 ..Config::default()
551 };
552 let m = resolve_model(&None, &config).unwrap().unwrap();
553 assert_eq!(m.provider, "openai");
554 assert_eq!(m.model, "gpt-sub");
555 }
556
557 #[test]
559 fn resolve_model_falls_back_to_defaults_chat() {
560 let config = Config {
561 defaults: Some(defaults_config(("anthropic", "claude-x"), None)),
562 ..Config::default()
563 };
564 let m = resolve_model(&None, &config).unwrap().unwrap();
565 assert_eq!(m.provider, "anthropic");
566 assert_eq!(m.model, "claude-x");
567 }
568
569 #[test]
572 fn resolve_model_none_when_nothing_configured() {
573 let config = Config {
574 defaults: None,
575 ..Config::default()
576 };
577 assert_eq!(resolve_model(&None, &config).unwrap(), None);
578 }
579
580 #[test]
583 fn resolve_model_malformed_colon_errors() {
584 let config = Config::default();
585 assert!(resolve_model(&some("openai:"), &config).is_err());
586 assert!(resolve_model(&some(":gpt-4o"), &config).is_err());
587 }
588}