1use std::collections::BTreeSet;
10use std::sync::Arc;
11
12use tokio::sync::{mpsc, RwLock};
13use tokio_util::sync::CancellationToken;
14
15use bamboo_agent_core::storage::{AttachmentReader, Storage};
16use bamboo_agent_core::tools::ToolExecutor;
17use bamboo_agent_core::{AgentEvent, Role, Session};
18use bamboo_config::PermissionMode;
19use bamboo_domain::ReasoningEffort;
20use bamboo_llm::Config;
21use bamboo_llm::LLMProvider;
22use bamboo_metrics::MetricsCollector;
23use bamboo_skills::SkillManager;
24
25use crate::runtime::config::{
26 AgentLoopConfig, AuxiliaryModelConfig, BashCompletionSink, BashResumeHook, GoldConfig,
27 GuardianConfig, GuardianSpawner, ImageFallbackConfig, PromptMemoryFlags,
28};
29use crate::runtime::hooks::HookRunner;
30use crate::runtime::model_roster::{ModelRoster, RoleModel};
31use crate::runtime::runner::run_agent_loop_with_config;
32use bamboo_domain::RuntimeSessionPersistence;
33
34#[derive(Clone)]
43pub struct AgentRuntime {
44 pub storage: Arc<dyn Storage>,
45 pub persistence: Arc<dyn RuntimeSessionPersistence>,
46 pub attachment_reader: Arc<dyn AttachmentReader>,
47 pub skill_manager: Arc<SkillManager>,
48 pub metrics_collector: MetricsCollector,
49 pub config: Arc<RwLock<Config>>,
50
51 pub provider: Arc<dyn LLMProvider>,
53
54 pub default_tools: Arc<dyn ToolExecutor>,
58
59 pub hook_runner: Arc<HookRunner>,
62}
63
64pub struct AgentRuntimeBuilder {
79 storage: Option<Arc<dyn Storage>>,
80 persistence: Option<Arc<dyn RuntimeSessionPersistence>>,
81 attachment_reader: Option<Arc<dyn AttachmentReader>>,
82 skill_manager: Option<Arc<SkillManager>>,
83 metrics_collector: Option<MetricsCollector>,
84 config: Option<Arc<RwLock<Config>>>,
85 provider: Option<Arc<dyn LLMProvider>>,
86 default_tools: Option<Arc<dyn ToolExecutor>>,
87 hook_runner: Arc<HookRunner>,
88}
89
90impl AgentRuntimeBuilder {
91 pub fn new() -> Self {
92 Self {
93 storage: None,
94 persistence: None,
95 attachment_reader: None,
96 skill_manager: None,
97 metrics_collector: None,
98 config: None,
99 provider: None,
100 default_tools: None,
101 hook_runner: Arc::new(HookRunner::new()),
102 }
103 }
104
105 pub fn storage(mut self, v: Arc<dyn Storage>) -> Self {
106 self.storage = Some(v);
107 self
108 }
109
110 pub fn persistence(mut self, v: Arc<dyn RuntimeSessionPersistence>) -> Self {
111 self.persistence = Some(v);
112 self
113 }
114
115 pub fn attachment_reader(mut self, v: Arc<dyn AttachmentReader>) -> Self {
116 self.attachment_reader = Some(v);
117 self
118 }
119
120 pub fn skill_manager(mut self, v: Arc<SkillManager>) -> Self {
121 self.skill_manager = Some(v);
122 self
123 }
124
125 pub fn metrics_collector(mut self, v: MetricsCollector) -> Self {
126 self.metrics_collector = Some(v);
127 self
128 }
129
130 pub fn config(mut self, v: Arc<RwLock<Config>>) -> Self {
131 self.config = Some(v);
132 self
133 }
134
135 pub fn provider(mut self, v: Arc<dyn LLMProvider>) -> Self {
136 self.provider = Some(v);
137 self
138 }
139
140 pub fn default_tools(mut self, v: Arc<dyn ToolExecutor>) -> Self {
141 self.default_tools = Some(v);
142 self
143 }
144
145 pub fn hook_runner(mut self, v: Arc<HookRunner>) -> Self {
147 self.hook_runner = v;
148 self
149 }
150
151 pub fn build(self) -> Result<AgentRuntime, &'static str> {
152 Ok(AgentRuntime {
153 storage: self.storage.ok_or_else(|| format_missing("storage"))?,
154 persistence: self
155 .persistence
156 .ok_or_else(|| format_missing("persistence"))?,
157 attachment_reader: self
158 .attachment_reader
159 .ok_or_else(|| format_missing("attachment_reader"))?,
160 skill_manager: self
161 .skill_manager
162 .ok_or_else(|| format_missing("skill_manager"))?,
163 metrics_collector: self
164 .metrics_collector
165 .ok_or_else(|| format_missing("metrics_collector"))?,
166 config: self.config.ok_or_else(|| format_missing("config"))?,
167 provider: self.provider.ok_or_else(|| format_missing("provider"))?,
168 default_tools: self
169 .default_tools
170 .ok_or_else(|| format_missing("default_tools"))?,
171 hook_runner: self.hook_runner,
172 })
173 }
174}
175
176fn format_missing(field: &str) -> &'static str {
177 match field {
180 "storage" => "AgentRuntimeBuilder: missing storage",
181 "persistence" => "AgentRuntimeBuilder: missing persistence",
182 "attachment_reader" => "AgentRuntimeBuilder: missing attachment_reader",
183 "skill_manager" => "AgentRuntimeBuilder: missing skill_manager",
184 "metrics_collector" => "AgentRuntimeBuilder: missing metrics_collector",
185 "config" => "AgentRuntimeBuilder: missing config",
186 "provider" => "AgentRuntimeBuilder: missing provider",
187 "default_tools" => "AgentRuntimeBuilder: missing default_tools",
188 _ => "AgentRuntimeBuilder: missing required field",
189 }
190}
191
192impl Default for AgentRuntimeBuilder {
193 fn default() -> Self {
194 Self::new()
195 }
196}
197
198pub struct ExecuteRequest {
208 pub initial_message: String,
210 pub event_tx: mpsc::Sender<AgentEvent>,
211 pub cancel_token: CancellationToken,
212
213 pub tools: Option<Arc<dyn ToolExecutor>>,
217 pub provider_override: Option<Arc<dyn LLMProvider>>,
220
221 pub model_roster: ModelRoster,
228 pub reasoning_effort: Option<ReasoningEffort>,
229 pub auxiliary_model_resolver: Option<Arc<dyn Fn() -> AuxiliaryModelConfig + Send + Sync>>,
232 pub disabled_filter_resolver:
235 Option<Arc<dyn Fn() -> (BTreeSet<String>, BTreeSet<String>) + Send + Sync>>,
236 pub disabled_tools: Option<BTreeSet<String>>,
238 pub disabled_skill_ids: Option<BTreeSet<String>>,
240 pub selected_skill_ids: Option<Vec<String>>,
241 pub selected_skill_mode: Option<String>,
242 pub image_fallback: Option<ImageFallbackConfig>,
243 pub gold_config: Option<GoldConfig>,
244 pub guardian_config: Option<GuardianConfig>,
246 pub guardian_spawner: Option<Arc<dyn GuardianSpawner>>,
249 pub bash_resume_hook: Option<Arc<dyn BashResumeHook>>,
252 pub bash_completion_sink: Option<Arc<dyn BashCompletionSink>>,
255 pub app_data_dir: Option<std::path::PathBuf>,
257 pub run_budget: Option<bamboo_config::RunBudgetConfig>,
265}
266
267pub struct ExecuteRequestBuilder {
283 initial_message: String,
284 event_tx: mpsc::Sender<AgentEvent>,
285 cancel_token: CancellationToken,
286
287 tools: Option<Arc<dyn ToolExecutor>>,
288 provider_override: Option<Arc<dyn LLMProvider>>,
289 model: Option<String>,
293 provider_name: Option<String>,
294 provider_type: Option<String>,
295 fast_model: Option<String>,
296 fast_model_provider: Option<Arc<dyn LLMProvider>>,
297 background_model: Option<String>,
298 background_model_provider: Option<Arc<dyn LLMProvider>>,
299 summarization_model: Option<String>,
300 summarization_model_provider: Option<Arc<dyn LLMProvider>>,
301 reasoning_effort: Option<ReasoningEffort>,
302 auxiliary_model_resolver: Option<Arc<dyn Fn() -> AuxiliaryModelConfig + Send + Sync>>,
303 disabled_filter_resolver:
304 Option<Arc<dyn Fn() -> (BTreeSet<String>, BTreeSet<String>) + Send + Sync>>,
305 disabled_tools: Option<BTreeSet<String>>,
306 disabled_skill_ids: Option<BTreeSet<String>>,
307 selected_skill_ids: Option<Vec<String>>,
308 selected_skill_mode: Option<String>,
309 image_fallback: Option<ImageFallbackConfig>,
310 gold_config: Option<GoldConfig>,
311 guardian_config: Option<GuardianConfig>,
312 guardian_spawner: Option<Arc<dyn GuardianSpawner>>,
313 bash_resume_hook: Option<Arc<dyn BashResumeHook>>,
314 bash_completion_sink: Option<Arc<dyn BashCompletionSink>>,
315 app_data_dir: Option<std::path::PathBuf>,
316 run_budget: Option<bamboo_config::RunBudgetConfig>,
317}
318
319impl ExecuteRequestBuilder {
320 pub fn new(
323 initial_message: impl Into<String>,
324 event_tx: mpsc::Sender<AgentEvent>,
325 cancel_token: CancellationToken,
326 ) -> Self {
327 Self {
328 initial_message: initial_message.into(),
329 event_tx,
330 cancel_token,
331 tools: None,
332 provider_override: None,
333 model: None,
334 provider_name: None,
335 provider_type: None,
336 fast_model: None,
337 fast_model_provider: None,
338 background_model: None,
339 background_model_provider: None,
340 summarization_model: None,
341 summarization_model_provider: None,
342 reasoning_effort: None,
343 auxiliary_model_resolver: None,
344 disabled_filter_resolver: None,
345 disabled_tools: None,
346 disabled_skill_ids: None,
347 selected_skill_ids: None,
348 selected_skill_mode: None,
349 image_fallback: None,
350 gold_config: None,
351 guardian_config: None,
352 guardian_spawner: None,
353 bash_resume_hook: None,
354 bash_completion_sink: None,
355 app_data_dir: None,
356 run_budget: None,
357 }
358 }
359
360 pub fn tools(mut self, v: Arc<dyn ToolExecutor>) -> Self {
362 self.tools = Some(v);
363 self
364 }
365
366 pub fn provider_override(mut self, v: Arc<dyn LLMProvider>) -> Self {
368 self.provider_override = Some(v);
369 self
370 }
371
372 pub fn model_roster(mut self, roster: ModelRoster) -> Self {
378 self.fast_model = roster.fast_model();
379 self.fast_model_provider = roster.fast_model_provider();
380 self.background_model = roster.background_model();
381 self.background_model_provider = roster.background_model_provider();
382 self.summarization_model = roster.summarization_model();
383 self.summarization_model_provider = roster.summarization_model_provider();
384 self.model = roster.model;
385 self.provider_name = roster.provider_name;
386 self.provider_type = roster.provider_type;
387 self
388 }
389
390 pub fn model(mut self, v: impl Into<String>) -> Self {
392 self.model = Some(v.into());
393 self
394 }
395
396 pub fn provider_name(mut self, v: impl Into<String>) -> Self {
398 self.provider_name = Some(v.into());
399 self
400 }
401
402 pub fn provider_type(mut self, v: impl Into<String>) -> Self {
404 self.provider_type = Some(v.into());
405 self
406 }
407
408 pub fn fast_model(mut self, v: impl Into<String>) -> Self {
410 self.fast_model = Some(v.into());
411 self
412 }
413
414 pub fn fast_model_provider(mut self, v: Arc<dyn LLMProvider>) -> Self {
416 self.fast_model_provider = Some(v);
417 self
418 }
419
420 pub fn background_model(mut self, v: impl Into<String>) -> Self {
422 self.background_model = Some(v.into());
423 self
424 }
425
426 pub fn background_model_provider(mut self, v: Arc<dyn LLMProvider>) -> Self {
428 self.background_model_provider = Some(v);
429 self
430 }
431
432 pub fn summarization_model(mut self, v: impl Into<String>) -> Self {
434 self.summarization_model = Some(v.into());
435 self
436 }
437
438 pub fn summarization_model_provider(mut self, v: Arc<dyn LLMProvider>) -> Self {
440 self.summarization_model_provider = Some(v);
441 self
442 }
443
444 pub fn reasoning_effort(mut self, v: ReasoningEffort) -> Self {
446 self.reasoning_effort = Some(v);
447 self
448 }
449
450 pub fn auxiliary_model_resolver(
452 mut self,
453 v: Arc<dyn Fn() -> AuxiliaryModelConfig + Send + Sync>,
454 ) -> Self {
455 self.auxiliary_model_resolver = Some(v);
456 self
457 }
458
459 pub fn disabled_filter_resolver(
461 mut self,
462 v: Arc<dyn Fn() -> (BTreeSet<String>, BTreeSet<String>) + Send + Sync>,
463 ) -> Self {
464 self.disabled_filter_resolver = Some(v);
465 self
466 }
467
468 pub fn disabled_tools(mut self, v: BTreeSet<String>) -> Self {
470 self.disabled_tools = Some(v);
471 self
472 }
473
474 pub fn disabled_skill_ids(mut self, v: BTreeSet<String>) -> Self {
476 self.disabled_skill_ids = Some(v);
477 self
478 }
479
480 pub fn selected_skill_ids(mut self, v: Vec<String>) -> Self {
482 self.selected_skill_ids = Some(v);
483 self
484 }
485
486 pub fn selected_skill_mode(mut self, v: impl Into<String>) -> Self {
488 self.selected_skill_mode = Some(v.into());
489 self
490 }
491
492 pub fn image_fallback(mut self, v: ImageFallbackConfig) -> Self {
494 self.image_fallback = Some(v);
495 self
496 }
497
498 pub(crate) fn gold_config(mut self, v: Option<GoldConfig>) -> Self {
505 self.gold_config = v;
506 self
507 }
508
509 pub(crate) fn guardian_config(mut self, v: Option<GuardianConfig>) -> Self {
512 self.guardian_config = v;
513 self
514 }
515
516 pub(crate) fn guardian_spawner(mut self, v: Option<Arc<dyn GuardianSpawner>>) -> Self {
519 self.guardian_spawner = v;
520 self
521 }
522
523 pub(crate) fn bash_resume_hook(mut self, v: Option<Arc<dyn BashResumeHook>>) -> Self {
526 self.bash_resume_hook = v;
527 self
528 }
529
530 pub(crate) fn bash_completion_sink(mut self, v: Option<Arc<dyn BashCompletionSink>>) -> Self {
533 self.bash_completion_sink = v;
534 self
535 }
536
537 pub fn app_data_dir(mut self, v: std::path::PathBuf) -> Self {
539 self.app_data_dir = Some(v);
540 self
541 }
542
543 pub fn run_budget(mut self, v: bamboo_config::RunBudgetConfig) -> Self {
548 self.run_budget = Some(v);
549 self
550 }
551
552 pub fn build(self) -> ExecuteRequest {
557 let model_roster = ModelRoster {
558 model: self.model,
559 provider_name: self.provider_name,
560 provider_type: self.provider_type,
561 fast: RoleModel::from_parts(self.fast_model, self.fast_model_provider),
562 background: RoleModel::from_parts(
563 self.background_model,
564 self.background_model_provider,
565 ),
566 summarization: RoleModel::from_parts(
567 self.summarization_model,
568 self.summarization_model_provider,
569 ),
570 };
571 ExecuteRequest {
572 initial_message: self.initial_message,
573 event_tx: self.event_tx,
574 cancel_token: self.cancel_token,
575 tools: self.tools,
576 provider_override: self.provider_override,
577 model_roster,
578 reasoning_effort: self.reasoning_effort,
579 auxiliary_model_resolver: self.auxiliary_model_resolver,
580 disabled_filter_resolver: self.disabled_filter_resolver,
581 disabled_tools: self.disabled_tools,
582 disabled_skill_ids: self.disabled_skill_ids,
583 selected_skill_ids: self.selected_skill_ids,
584 selected_skill_mode: self.selected_skill_mode,
585 image_fallback: self.image_fallback,
586 gold_config: self.gold_config,
587 guardian_config: self.guardian_config,
588 guardian_spawner: self.guardian_spawner,
589 bash_resume_hook: self.bash_resume_hook,
590 bash_completion_sink: self.bash_completion_sink,
591 app_data_dir: self.app_data_dir,
592 run_budget: self.run_budget,
593 }
594 }
595}
596
597fn extract_system_prompt(session: &Session) -> Option<String> {
603 session
604 .messages
605 .iter()
606 .find(|m| matches!(m.role, Role::System))
607 .map(|m| m.content.clone())
608}
609
610impl AgentRuntime {
615 pub async fn execute(
620 &self,
621 session: &mut Session,
622 req: ExecuteRequest,
623 ) -> crate::runtime::runner::Result<()> {
624 let system_prompt = extract_system_prompt(session);
625 let config = self.config.read().await;
626 let ExecuteRequest {
627 initial_message,
628 event_tx,
629 cancel_token,
630 tools,
631 provider_override,
632 model_roster,
633 reasoning_effort,
634 auxiliary_model_resolver,
635 disabled_filter_resolver,
636 disabled_tools,
637 disabled_skill_ids,
638 selected_skill_ids,
639 selected_skill_mode,
640 image_fallback,
641 gold_config,
642 guardian_config,
643 guardian_spawner,
644 bash_resume_hook,
645 bash_completion_sink,
646 app_data_dir,
647 run_budget,
648 } = req;
649 let tools = tools.unwrap_or_else(|| self.default_tools.clone());
650 let llm = provider_override.unwrap_or_else(|| self.provider.clone());
651
652 let fast_model = model_roster.fast_model();
657 let fast_model_provider = model_roster.fast_model_provider();
658 let background_model = model_roster.background_model();
659 let background_model_provider = model_roster.background_model_provider();
660 let summarization_model = model_roster.summarization_model();
661 let summarization_model_provider = model_roster.summarization_model_provider();
662 let ModelRoster {
663 model,
664 provider_name,
665 provider_type,
666 ..
667 } = model_roster;
668 let hook_runner = Arc::new(
669 self.hook_runner
670 .with_lifecycle_config(&config.lifecycle_hooks, app_data_dir.clone()),
671 );
672
673 let loop_config = AgentLoopConfig {
674 max_rounds: 200,
675 system_prompt,
676 legacy_model_limits: config.extra.get("model_limits").cloned(),
679 disabled_skill_ids: disabled_skill_ids.unwrap_or_else(|| config.disabled_skill_ids()),
680 selected_skill_ids,
681 selected_skill_mode,
682 skill_manager: Some(self.skill_manager.clone()),
683 skip_initial_user_message: true,
684 storage: Some(self.storage.clone()),
685 persistence: Some(self.persistence.clone()),
686 attachment_reader: Some(self.attachment_reader.clone()),
687 metrics_collector: Some(self.metrics_collector.clone()),
688 model_name: model,
689 fast_model_name: fast_model.or_else(|| config.get_fast_model()),
690 fast_model_provider,
691 background_model_name: background_model
692 .or_else(|| config.get_memory_background_model()),
693 planning_model_name: config
694 .defaults
695 .as_ref()
696 .and_then(|d| d.planning.as_ref())
697 .map(|r| r.model.clone()),
698 search_model_name: config
699 .defaults
700 .as_ref()
701 .and_then(|d| d.search.as_ref().or(d.fast.as_ref()))
702 .map(|r| r.model.clone()),
703 compression_instructions: None,
704 summarization_model_name: summarization_model
705 .or_else(|| config.get_task_summary_model()),
706 background_model_provider,
707 summarization_model_provider,
708 provider_name: Some(provider_name.unwrap_or_else(|| config.provider.clone())),
709 provider_type,
710 reasoning_effort,
711 auxiliary_model_resolver,
712 disabled_filter_resolver,
713 disabled_tools: {
714 let mut merged = config.disabled_tool_names();
715 if let Some(dt) = disabled_tools {
716 merged.extend(dt);
717 }
718 merged
719 },
720 image_fallback,
721 app_data_dir,
722 prompt_memory_flags: config
723 .memory()
724 .as_ref()
725 .map(PromptMemoryFlags::from)
726 .unwrap_or_default(),
727 features_dynamic_model_routing: config.features.dynamic_model_routing,
728 permission_mode: session
729 .agent_runtime_state
730 .as_ref()
731 .and_then(|state| state.plan_mode.as_ref())
732 .map(|_| PermissionMode::Plan),
733 gold_config,
734 guardian_config,
735 guardian_spawner,
736 bash_resume_hook,
737 bash_completion_sink,
738 hook_runner,
739 mcp_tool_guidance: tools.tool_guidance(),
743 run_budget: config.run_budget.merged_with_override(run_budget.as_ref()),
749 stream_timeout: config.stream_timeout,
750 ..Default::default()
751 };
752
753 drop(config);
754
755 let session_end_runner = loop_config.hook_runner.clone();
756 let session_end_event_tx = event_tx.clone();
757 let result = run_agent_loop_with_config(
758 session,
759 initial_message,
760 event_tx,
761 llm,
762 tools,
763 cancel_token,
764 loop_config,
765 )
766 .await;
767
768 crate::runtime::hooks::run_session_end_hooks(
769 &session_end_runner,
770 &result,
771 session,
772 &session_end_event_tx,
773 )
774 .await;
775
776 if let Err(checkpoint_error) = self.persistence.checkpoint_runtime_session(session).await {
788 match &result {
789 Ok(()) => tracing::warn!(
790 session_id = %session.id,
791 error = %checkpoint_error,
792 "failed to checkpoint session transcript after successful execution"
793 ),
794 Err(execution_error) => tracing::warn!(
795 session_id = %session.id,
796 error = %checkpoint_error,
797 execution_error = %execution_error,
798 "failed to checkpoint session transcript after execution error"
799 ),
800 }
801 }
802
803 result
804 }
805}