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, SessionInboxPort};
33
34use crate::session_activation::SessionActivationRouter;
35use crate::session_messaging::SessionMessenger;
36
37#[derive(Clone)]
46pub struct AgentRuntime {
47 pub storage: Arc<dyn Storage>,
48 pub persistence: Arc<dyn RuntimeSessionPersistence>,
49 pub session_inbox: Option<Arc<dyn SessionInboxPort>>,
50 pub activation_router: Option<Arc<SessionActivationRouter>>,
51 pub session_messenger: Option<Arc<SessionMessenger>>,
52 pub attachment_reader: Arc<dyn AttachmentReader>,
53 pub skill_manager: Arc<SkillManager>,
54 pub project_context_resolver: Option<Arc<crate::project_context::ProjectContextResolver>>,
55 pub metrics_collector: MetricsCollector,
56 pub config: Arc<RwLock<Config>>,
57
58 pub provider: Arc<dyn LLMProvider>,
60
61 pub default_tools: Arc<dyn ToolExecutor>,
65
66 pub hook_runner: Arc<HookRunner>,
69}
70
71pub struct AgentRuntimeBuilder {
86 storage: Option<Arc<dyn Storage>>,
87 persistence: Option<Arc<dyn RuntimeSessionPersistence>>,
88 session_inbox: Option<Arc<dyn SessionInboxPort>>,
89 activation_router: Option<Arc<SessionActivationRouter>>,
90 session_messenger: Option<Arc<SessionMessenger>>,
91 attachment_reader: Option<Arc<dyn AttachmentReader>>,
92 skill_manager: Option<Arc<SkillManager>>,
93 project_context_resolver: Option<Arc<crate::project_context::ProjectContextResolver>>,
94 metrics_collector: Option<MetricsCollector>,
95 config: Option<Arc<RwLock<Config>>>,
96 provider: Option<Arc<dyn LLMProvider>>,
97 default_tools: Option<Arc<dyn ToolExecutor>>,
98 hook_runner: Arc<HookRunner>,
99}
100
101impl AgentRuntimeBuilder {
102 pub fn new() -> Self {
103 Self {
104 storage: None,
105 persistence: None,
106 session_inbox: None,
107 activation_router: None,
108 session_messenger: None,
109 attachment_reader: None,
110 skill_manager: None,
111 project_context_resolver: None,
112 metrics_collector: None,
113 config: None,
114 provider: None,
115 default_tools: None,
116 hook_runner: Arc::new(HookRunner::new()),
117 }
118 }
119
120 pub fn storage(mut self, v: Arc<dyn Storage>) -> Self {
121 self.storage = Some(v);
122 self
123 }
124
125 pub fn persistence(mut self, v: Arc<dyn RuntimeSessionPersistence>) -> Self {
126 self.persistence = Some(v);
127 self
128 }
129
130 pub fn session_inbox(mut self, v: Arc<dyn SessionInboxPort>) -> Self {
131 self.session_inbox = Some(v);
132 self
133 }
134
135 pub fn activation_router(mut self, v: Arc<SessionActivationRouter>) -> Self {
136 self.activation_router = Some(v);
137 self
138 }
139
140 pub fn session_messenger(mut self, v: Arc<SessionMessenger>) -> Self {
141 self.session_messenger = Some(v);
142 self
143 }
144
145 pub fn attachment_reader(mut self, v: Arc<dyn AttachmentReader>) -> Self {
146 self.attachment_reader = Some(v);
147 self
148 }
149
150 pub fn skill_manager(mut self, v: Arc<SkillManager>) -> Self {
151 self.skill_manager = Some(v);
152 self
153 }
154
155 pub fn project_context_resolver(
156 mut self,
157 v: Arc<crate::project_context::ProjectContextResolver>,
158 ) -> Self {
159 self.project_context_resolver = Some(v);
160 self
161 }
162
163 pub fn metrics_collector(mut self, v: MetricsCollector) -> Self {
164 self.metrics_collector = Some(v);
165 self
166 }
167
168 pub fn config(mut self, v: Arc<RwLock<Config>>) -> Self {
169 self.config = Some(v);
170 self
171 }
172
173 pub fn provider(mut self, v: Arc<dyn LLMProvider>) -> Self {
174 self.provider = Some(v);
175 self
176 }
177
178 pub fn default_tools(mut self, v: Arc<dyn ToolExecutor>) -> Self {
179 self.default_tools = Some(v);
180 self
181 }
182
183 pub fn hook_runner(mut self, v: Arc<HookRunner>) -> Self {
185 self.hook_runner = v;
186 self
187 }
188
189 pub fn build(self) -> Result<AgentRuntime, &'static str> {
190 if let (Some(router), Some(inbox)) = (&self.activation_router, &self.session_inbox) {
191 router.set_inbox(inbox.clone());
192 }
193 Ok(AgentRuntime {
194 storage: self.storage.ok_or_else(|| format_missing("storage"))?,
195 persistence: self
196 .persistence
197 .ok_or_else(|| format_missing("persistence"))?,
198 session_inbox: self.session_inbox,
199 activation_router: self.activation_router,
200 session_messenger: self.session_messenger,
201 attachment_reader: self
202 .attachment_reader
203 .ok_or_else(|| format_missing("attachment_reader"))?,
204 skill_manager: self
205 .skill_manager
206 .ok_or_else(|| format_missing("skill_manager"))?,
207 project_context_resolver: self.project_context_resolver,
208 metrics_collector: self
209 .metrics_collector
210 .ok_or_else(|| format_missing("metrics_collector"))?,
211 config: self.config.ok_or_else(|| format_missing("config"))?,
212 provider: self.provider.ok_or_else(|| format_missing("provider"))?,
213 default_tools: self
214 .default_tools
215 .ok_or_else(|| format_missing("default_tools"))?,
216 hook_runner: self.hook_runner,
217 })
218 }
219}
220
221fn format_missing(field: &str) -> &'static str {
222 match field {
225 "storage" => "AgentRuntimeBuilder: missing storage",
226 "persistence" => "AgentRuntimeBuilder: missing persistence",
227 "attachment_reader" => "AgentRuntimeBuilder: missing attachment_reader",
228 "skill_manager" => "AgentRuntimeBuilder: missing skill_manager",
229 "metrics_collector" => "AgentRuntimeBuilder: missing metrics_collector",
230 "config" => "AgentRuntimeBuilder: missing config",
231 "provider" => "AgentRuntimeBuilder: missing provider",
232 "default_tools" => "AgentRuntimeBuilder: missing default_tools",
233 _ => "AgentRuntimeBuilder: missing required field",
234 }
235}
236
237impl Default for AgentRuntimeBuilder {
238 fn default() -> Self {
239 Self::new()
240 }
241}
242
243pub struct ExecuteRequest {
253 pub initial_message: String,
255 pub event_tx: mpsc::Sender<AgentEvent>,
256 pub cancel_token: CancellationToken,
257
258 pub tools: Option<Arc<dyn ToolExecutor>>,
262 pub provider_override: Option<Arc<dyn LLMProvider>>,
265
266 pub model_roster: ModelRoster,
273 pub reasoning_effort: Option<ReasoningEffort>,
274 pub auxiliary_model_resolver: Option<Arc<dyn Fn() -> AuxiliaryModelConfig + Send + Sync>>,
277 pub disabled_filter_resolver:
280 Option<Arc<dyn Fn() -> (BTreeSet<String>, BTreeSet<String>) + Send + Sync>>,
281 pub disabled_tools: Option<BTreeSet<String>>,
283 pub disabled_skill_ids: Option<BTreeSet<String>>,
285 pub selected_skill_ids: Option<Vec<String>>,
286 pub selected_skill_mode: Option<String>,
287 pub image_fallback: Option<ImageFallbackConfig>,
288 pub gold_config: Option<GoldConfig>,
289 pub guardian_config: Option<GuardianConfig>,
291 pub guardian_spawner: Option<Arc<dyn GuardianSpawner>>,
294 pub bash_resume_hook: Option<Arc<dyn BashResumeHook>>,
297 pub bash_completion_sink: Option<Arc<dyn BashCompletionSink>>,
300 pub app_data_dir: Option<std::path::PathBuf>,
302 pub run_budget: Option<bamboo_config::RunBudgetConfig>,
310}
311
312pub struct ExecuteRequestBuilder {
328 initial_message: String,
329 event_tx: mpsc::Sender<AgentEvent>,
330 cancel_token: CancellationToken,
331
332 tools: Option<Arc<dyn ToolExecutor>>,
333 provider_override: Option<Arc<dyn LLMProvider>>,
334 model: Option<String>,
338 provider_name: Option<String>,
339 provider_type: Option<String>,
340 fast_model: Option<String>,
341 fast_model_provider: Option<Arc<dyn LLMProvider>>,
342 background_model: Option<String>,
343 background_model_provider: Option<Arc<dyn LLMProvider>>,
344 summarization_model: Option<String>,
345 summarization_model_provider: Option<Arc<dyn LLMProvider>>,
346 reasoning_effort: Option<ReasoningEffort>,
347 auxiliary_model_resolver: Option<Arc<dyn Fn() -> AuxiliaryModelConfig + Send + Sync>>,
348 disabled_filter_resolver:
349 Option<Arc<dyn Fn() -> (BTreeSet<String>, BTreeSet<String>) + Send + Sync>>,
350 disabled_tools: Option<BTreeSet<String>>,
351 disabled_skill_ids: Option<BTreeSet<String>>,
352 selected_skill_ids: Option<Vec<String>>,
353 selected_skill_mode: Option<String>,
354 image_fallback: Option<ImageFallbackConfig>,
355 gold_config: Option<GoldConfig>,
356 guardian_config: Option<GuardianConfig>,
357 guardian_spawner: Option<Arc<dyn GuardianSpawner>>,
358 bash_resume_hook: Option<Arc<dyn BashResumeHook>>,
359 bash_completion_sink: Option<Arc<dyn BashCompletionSink>>,
360 app_data_dir: Option<std::path::PathBuf>,
361 run_budget: Option<bamboo_config::RunBudgetConfig>,
362}
363
364impl ExecuteRequestBuilder {
365 pub fn new(
368 initial_message: impl Into<String>,
369 event_tx: mpsc::Sender<AgentEvent>,
370 cancel_token: CancellationToken,
371 ) -> Self {
372 Self {
373 initial_message: initial_message.into(),
374 event_tx,
375 cancel_token,
376 tools: None,
377 provider_override: None,
378 model: None,
379 provider_name: None,
380 provider_type: None,
381 fast_model: None,
382 fast_model_provider: None,
383 background_model: None,
384 background_model_provider: None,
385 summarization_model: None,
386 summarization_model_provider: None,
387 reasoning_effort: None,
388 auxiliary_model_resolver: None,
389 disabled_filter_resolver: None,
390 disabled_tools: None,
391 disabled_skill_ids: None,
392 selected_skill_ids: None,
393 selected_skill_mode: None,
394 image_fallback: None,
395 gold_config: None,
396 guardian_config: None,
397 guardian_spawner: None,
398 bash_resume_hook: None,
399 bash_completion_sink: None,
400 app_data_dir: None,
401 run_budget: None,
402 }
403 }
404
405 pub fn tools(mut self, v: Arc<dyn ToolExecutor>) -> Self {
407 self.tools = Some(v);
408 self
409 }
410
411 pub fn provider_override(mut self, v: Arc<dyn LLMProvider>) -> Self {
413 self.provider_override = Some(v);
414 self
415 }
416
417 pub fn model_roster(mut self, roster: ModelRoster) -> Self {
423 self.fast_model = roster.fast_model();
424 self.fast_model_provider = roster.fast_model_provider();
425 self.background_model = roster.background_model();
426 self.background_model_provider = roster.background_model_provider();
427 self.summarization_model = roster.summarization_model();
428 self.summarization_model_provider = roster.summarization_model_provider();
429 self.model = roster.model;
430 self.provider_name = roster.provider_name;
431 self.provider_type = roster.provider_type;
432 self
433 }
434
435 pub fn model(mut self, v: impl Into<String>) -> Self {
437 self.model = Some(v.into());
438 self
439 }
440
441 pub fn provider_name(mut self, v: impl Into<String>) -> Self {
443 self.provider_name = Some(v.into());
444 self
445 }
446
447 pub fn provider_type(mut self, v: impl Into<String>) -> Self {
449 self.provider_type = Some(v.into());
450 self
451 }
452
453 pub fn fast_model(mut self, v: impl Into<String>) -> Self {
455 self.fast_model = Some(v.into());
456 self
457 }
458
459 pub fn fast_model_provider(mut self, v: Arc<dyn LLMProvider>) -> Self {
461 self.fast_model_provider = Some(v);
462 self
463 }
464
465 pub fn background_model(mut self, v: impl Into<String>) -> Self {
467 self.background_model = Some(v.into());
468 self
469 }
470
471 pub fn background_model_provider(mut self, v: Arc<dyn LLMProvider>) -> Self {
473 self.background_model_provider = Some(v);
474 self
475 }
476
477 pub fn summarization_model(mut self, v: impl Into<String>) -> Self {
479 self.summarization_model = Some(v.into());
480 self
481 }
482
483 pub fn summarization_model_provider(mut self, v: Arc<dyn LLMProvider>) -> Self {
485 self.summarization_model_provider = Some(v);
486 self
487 }
488
489 pub fn reasoning_effort(mut self, v: ReasoningEffort) -> Self {
491 self.reasoning_effort = Some(v);
492 self
493 }
494
495 pub fn auxiliary_model_resolver(
497 mut self,
498 v: Arc<dyn Fn() -> AuxiliaryModelConfig + Send + Sync>,
499 ) -> Self {
500 self.auxiliary_model_resolver = Some(v);
501 self
502 }
503
504 pub fn disabled_filter_resolver(
506 mut self,
507 v: Arc<dyn Fn() -> (BTreeSet<String>, BTreeSet<String>) + Send + Sync>,
508 ) -> Self {
509 self.disabled_filter_resolver = Some(v);
510 self
511 }
512
513 pub fn disabled_tools(mut self, v: BTreeSet<String>) -> Self {
515 self.disabled_tools = Some(v);
516 self
517 }
518
519 pub fn disabled_skill_ids(mut self, v: BTreeSet<String>) -> Self {
521 self.disabled_skill_ids = Some(v);
522 self
523 }
524
525 pub fn selected_skill_ids(mut self, v: Vec<String>) -> Self {
527 self.selected_skill_ids = Some(v);
528 self
529 }
530
531 pub fn selected_skill_mode(mut self, v: impl Into<String>) -> Self {
533 self.selected_skill_mode = Some(v.into());
534 self
535 }
536
537 pub fn image_fallback(mut self, v: ImageFallbackConfig) -> Self {
539 self.image_fallback = Some(v);
540 self
541 }
542
543 pub(crate) fn gold_config(mut self, v: Option<GoldConfig>) -> Self {
550 self.gold_config = v;
551 self
552 }
553
554 pub(crate) fn guardian_config(mut self, v: Option<GuardianConfig>) -> Self {
557 self.guardian_config = v;
558 self
559 }
560
561 pub(crate) fn guardian_spawner(mut self, v: Option<Arc<dyn GuardianSpawner>>) -> Self {
564 self.guardian_spawner = v;
565 self
566 }
567
568 pub(crate) fn bash_resume_hook(mut self, v: Option<Arc<dyn BashResumeHook>>) -> Self {
571 self.bash_resume_hook = v;
572 self
573 }
574
575 pub(crate) fn bash_completion_sink(mut self, v: Option<Arc<dyn BashCompletionSink>>) -> Self {
578 self.bash_completion_sink = v;
579 self
580 }
581
582 pub fn app_data_dir(mut self, v: std::path::PathBuf) -> Self {
584 self.app_data_dir = Some(v);
585 self
586 }
587
588 pub fn run_budget(mut self, v: bamboo_config::RunBudgetConfig) -> Self {
593 self.run_budget = Some(v);
594 self
595 }
596
597 pub fn build(self) -> ExecuteRequest {
602 let model_roster = ModelRoster {
603 model: self.model,
604 provider_name: self.provider_name,
605 provider_type: self.provider_type,
606 fast: RoleModel::from_parts(self.fast_model, self.fast_model_provider),
607 background: RoleModel::from_parts(
608 self.background_model,
609 self.background_model_provider,
610 ),
611 summarization: RoleModel::from_parts(
612 self.summarization_model,
613 self.summarization_model_provider,
614 ),
615 };
616 ExecuteRequest {
617 initial_message: self.initial_message,
618 event_tx: self.event_tx,
619 cancel_token: self.cancel_token,
620 tools: self.tools,
621 provider_override: self.provider_override,
622 model_roster,
623 reasoning_effort: self.reasoning_effort,
624 auxiliary_model_resolver: self.auxiliary_model_resolver,
625 disabled_filter_resolver: self.disabled_filter_resolver,
626 disabled_tools: self.disabled_tools,
627 disabled_skill_ids: self.disabled_skill_ids,
628 selected_skill_ids: self.selected_skill_ids,
629 selected_skill_mode: self.selected_skill_mode,
630 image_fallback: self.image_fallback,
631 gold_config: self.gold_config,
632 guardian_config: self.guardian_config,
633 guardian_spawner: self.guardian_spawner,
634 bash_resume_hook: self.bash_resume_hook,
635 bash_completion_sink: self.bash_completion_sink,
636 app_data_dir: self.app_data_dir,
637 run_budget: self.run_budget,
638 }
639 }
640}
641
642fn extract_system_prompt(session: &Session) -> Option<String> {
648 session
649 .messages
650 .iter()
651 .find(|m| matches!(m.role, Role::System))
652 .map(|m| m.content.clone())
653}
654
655impl AgentRuntime {
660 pub async fn execute(
665 &self,
666 session: &mut Session,
667 req: ExecuteRequest,
668 ) -> crate::runtime::runner::Result<()> {
669 let session_activation_notifications = match self.activation_router.as_ref() {
670 Some(router) => Some(Arc::new(parking_lot::Mutex::new(
671 router.subscribe(&session.id).await,
672 ))),
673 None => None,
674 };
675 let system_prompt = extract_system_prompt(session);
676 let config = self.config.read().await;
677 let ExecuteRequest {
678 initial_message,
679 event_tx,
680 cancel_token,
681 tools,
682 provider_override,
683 model_roster,
684 reasoning_effort,
685 auxiliary_model_resolver,
686 disabled_filter_resolver,
687 disabled_tools,
688 disabled_skill_ids,
689 selected_skill_ids,
690 selected_skill_mode,
691 image_fallback,
692 gold_config,
693 guardian_config,
694 guardian_spawner,
695 bash_resume_hook,
696 bash_completion_sink,
697 app_data_dir,
698 run_budget,
699 } = req;
700 let tools = tools.unwrap_or_else(|| self.default_tools.clone());
701 let llm = provider_override.unwrap_or_else(|| self.provider.clone());
702
703 let fast_model = model_roster.fast_model();
708 let fast_model_provider = model_roster.fast_model_provider();
709 let background_model = model_roster.background_model();
710 let background_model_provider = model_roster.background_model_provider();
711 let summarization_model = model_roster.summarization_model();
712 let summarization_model_provider = model_roster.summarization_model_provider();
713 let ModelRoster {
714 model,
715 provider_name,
716 provider_type,
717 ..
718 } = model_roster;
719 let hook_runner = Arc::new(
720 self.hook_runner
721 .with_lifecycle_config(&config.lifecycle_hooks, app_data_dir.clone()),
722 );
723
724 let loop_config = AgentLoopConfig {
725 max_rounds: 200,
726 system_prompt,
727 legacy_model_limits: config.extra.get("model_limits").cloned(),
730 disabled_skill_ids: disabled_skill_ids.unwrap_or_else(|| config.disabled_skill_ids()),
731 selected_skill_ids,
732 selected_skill_mode,
733 skill_manager: Some(self.skill_manager.clone()),
734 project_context_resolver: self.project_context_resolver.clone(),
735 skip_initial_user_message: true,
736 storage: Some(self.storage.clone()),
737 persistence: Some(self.persistence.clone()),
738 session_inbox: self.session_inbox.clone(),
739 session_activation_notifications,
740 attachment_reader: Some(self.attachment_reader.clone()),
741 metrics_collector: Some(self.metrics_collector.clone()),
742 model_name: model,
743 fast_model_name: fast_model.or_else(|| config.get_fast_model()),
744 fast_model_provider,
745 background_model_name: background_model
746 .or_else(|| config.get_memory_background_model()),
747 planning_model_name: config
748 .defaults
749 .as_ref()
750 .and_then(|d| d.planning.as_ref())
751 .map(|r| r.model.clone()),
752 search_model_name: config
753 .defaults
754 .as_ref()
755 .and_then(|d| d.search.as_ref().or(d.fast.as_ref()))
756 .map(|r| r.model.clone()),
757 compression_instructions: None,
758 summarization_model_name: summarization_model
759 .or_else(|| config.get_task_summary_model()),
760 background_model_provider,
761 summarization_model_provider,
762 provider_name: Some(provider_name.unwrap_or_else(|| config.provider.clone())),
763 provider_type,
764 reasoning_effort,
765 auxiliary_model_resolver,
766 disabled_filter_resolver,
767 disabled_tools: {
768 let mut merged = config.disabled_tool_names();
769 if let Some(dt) = disabled_tools {
770 merged.extend(dt);
771 }
772 merged
773 },
774 image_fallback,
775 app_data_dir,
776 prompt_memory_flags: config
777 .memory()
778 .as_ref()
779 .map(PromptMemoryFlags::from)
780 .unwrap_or_default(),
781 features_dynamic_model_routing: config.features.dynamic_model_routing,
782 permission_mode: session
783 .agent_runtime_state
784 .as_ref()
785 .and_then(|state| state.plan_mode.as_ref())
786 .map(|_| PermissionMode::Plan),
787 gold_config,
788 guardian_config,
789 guardian_spawner,
790 bash_resume_hook,
791 bash_completion_sink,
792 hook_runner,
793 mcp_tool_guidance: tools.tool_guidance(),
797 run_budget: config.run_budget.merged_with_override(run_budget.as_ref()),
803 stream_timeout: config.stream_timeout,
804 ..Default::default()
805 };
806
807 drop(config);
808
809 let session_end_runner = loop_config.hook_runner.clone();
810 let session_end_event_tx = event_tx.clone();
811 let result = run_agent_loop_with_config(
812 session,
813 initial_message,
814 event_tx,
815 llm,
816 tools,
817 cancel_token,
818 loop_config,
819 )
820 .await;
821
822 crate::runtime::hooks::run_session_end_hooks(
823 &session_end_runner,
824 &result,
825 session,
826 &session_end_event_tx,
827 )
828 .await;
829
830 if let Err(checkpoint_error) = self.persistence.checkpoint_runtime_session(session).await {
842 match &result {
843 Ok(()) => tracing::warn!(
844 session_id = %session.id,
845 error = %checkpoint_error,
846 "failed to checkpoint session transcript after successful execution"
847 ),
848 Err(execution_error) => tracing::warn!(
849 session_id = %session.id,
850 error = %checkpoint_error,
851 execution_error = %execution_error,
852 "failed to checkpoint session transcript after execution error"
853 ),
854 }
855 }
856
857 result
858 }
859}