1use std::collections::{BTreeSet, HashMap};
7use std::future::Future;
8use std::pin::Pin;
9use std::sync::Arc;
10
11use tokio::sync::{mpsc, RwLock};
12use tokio_util::sync::CancellationToken;
13use tracing::Instrument;
14
15use bamboo_agent_core::tools::ToolExecutor;
16use bamboo_agent_core::{AgentError, AgentEvent, Session};
17use bamboo_domain::ReasoningEffort;
18use bamboo_llm::LLMProvider;
19
20use crate::runtime::config::{
21 AuxiliaryModelConfig, BashCompletionSink, BashResumeHook, GoldConfig, GuardianConfig,
22 GuardianSpawner, ImageFallbackConfig,
23};
24use crate::runtime::execution::child_completion::ChildCompletion;
25use crate::runtime::execution::runner_lifecycle::finalize_runner;
26use crate::runtime::execution::runner_state::AgentRunner;
27use crate::runtime::model_roster::ModelRoster;
28use crate::runtime::Agent;
29use crate::runtime::{ExecuteRequest, ExecuteRequestBuilder};
30
31pub type SessionCache = std::sync::Arc<
39 dashmap::DashMap<String, std::sync::Arc<parking_lot::RwLock<bamboo_agent_core::Session>>>,
40>;
41
42pub fn read_cached_session(cache: &SessionCache, id: &str) -> Option<bamboo_agent_core::Session> {
50 cache
51 .get(id)
52 .map(|e| e.value().clone())
53 .map(|a| a.read().clone())
54}
55
56const SKILL_CONTEXT_START_MARKER: &str = "<!-- BAMBOO_SKILL_CONTEXT_START -->";
57const TOOL_GUIDE_START_MARKER: &str = "<!-- BAMBOO_TOOL_GUIDE_START -->";
58const EXTERNAL_MEMORY_START_MARKER: &str = "<!-- BAMBOO_EXTERNAL_MEMORY_START -->";
59const TASK_LIST_START_MARKER: &str = "<!-- BAMBOO_TASK_LIST_START -->";
60
61pub struct SessionExecutionOutcome {
68 pub success: bool,
70 pub cancelled: bool,
72 pub error: Option<String>,
74}
75
76impl SessionExecutionOutcome {
77 fn from_result(result: &Result<(), AgentError>) -> Self {
78 match result {
79 Ok(()) => Self {
80 success: true,
81 cancelled: false,
82 error: None,
83 },
84 Err(error) => Self {
85 success: false,
86 cancelled: error.is_cancelled(),
87 error: Some(error.to_string()),
88 },
89 }
90 }
91}
92
93pub type SessionCompletionHook = Box<
103 dyn for<'a> FnOnce(
104 SessionExecutionOutcome,
105 &'a mut Session,
106 ) -> Pin<Box<dyn Future<Output = ()> + Send + 'a>>
107 + Send,
108>;
109
110pub struct SessionExecutionArgs {
116 pub agent: Arc<Agent>,
118 pub session_id: String,
119 pub session: Session,
120
121 pub tools_override: Option<Arc<dyn ToolExecutor>>,
123 pub provider_override: Option<Arc<dyn LLMProvider>>,
124 pub model_roster: ModelRoster,
128 pub reasoning_effort: Option<ReasoningEffort>,
129 pub reasoning_effort_source: String,
130 pub auxiliary_model_resolver:
131 Option<Arc<dyn Fn() -> crate::runtime::config::AuxiliaryModelConfig + Send + Sync>>,
132 pub disabled_filter_resolver:
136 Option<Arc<dyn Fn() -> (BTreeSet<String>, BTreeSet<String>) + Send + Sync>>,
137 pub disabled_tools: Option<BTreeSet<String>>,
138 pub disabled_skill_ids: Option<BTreeSet<String>>,
139 pub selected_skill_ids: Option<Vec<String>>,
140 pub selected_skill_mode: Option<String>,
141 pub cancel_token: CancellationToken,
142 pub mpsc_tx: mpsc::Sender<AgentEvent>,
143 pub image_fallback: Option<ImageFallbackConfig>,
144 pub gold_config: Option<GoldConfig>,
145 pub guardian_config: Option<GuardianConfig>,
147 pub guardian_spawner: Option<Arc<dyn GuardianSpawner>>,
150 pub bash_resume_hook: Option<Arc<dyn BashResumeHook>>,
152 pub bash_completion_sink: Option<Arc<dyn BashCompletionSink>>,
154 pub app_data_dir: Option<std::path::PathBuf>,
155 pub run_budget: Option<bamboo_config::RunBudgetConfig>,
158
159 pub runners: Arc<RwLock<HashMap<String, AgentRunner>>>,
161 pub sessions_cache: SessionCache,
162
163 pub on_complete: Option<SessionCompletionHook>,
166
167 pub child_completion_handler: Option<Arc<dyn super::ChildCompletionHandler>>,
179}
180
181struct ExecuteRequestParams {
191 tools: Option<Arc<dyn ToolExecutor>>,
192 provider_override: Option<Arc<dyn LLMProvider>>,
193 model_roster: ModelRoster,
194 reasoning_effort: Option<ReasoningEffort>,
195 auxiliary_model_resolver: Option<Arc<dyn Fn() -> AuxiliaryModelConfig + Send + Sync>>,
196 disabled_filter_resolver:
197 Option<Arc<dyn Fn() -> (BTreeSet<String>, BTreeSet<String>) + Send + Sync>>,
198 disabled_tools: Option<BTreeSet<String>>,
199 disabled_skill_ids: Option<BTreeSet<String>>,
200 selected_skill_ids: Option<Vec<String>>,
201 selected_skill_mode: Option<String>,
202 image_fallback: Option<ImageFallbackConfig>,
203 gold_config: Option<GoldConfig>,
204 guardian_config: Option<GuardianConfig>,
205 guardian_spawner: Option<Arc<dyn GuardianSpawner>>,
206 bash_resume_hook: Option<Arc<dyn BashResumeHook>>,
207 bash_completion_sink: Option<Arc<dyn BashCompletionSink>>,
208 app_data_dir: Option<std::path::PathBuf>,
209 run_budget: Option<bamboo_config::RunBudgetConfig>,
210}
211
212fn build_execute_request(
219 initial_message: String,
220 event_tx: mpsc::Sender<AgentEvent>,
221 cancel_token: CancellationToken,
222 params: ExecuteRequestParams,
223) -> ExecuteRequest {
224 let ExecuteRequestParams {
225 tools,
226 provider_override,
227 model_roster,
228 reasoning_effort,
229 auxiliary_model_resolver,
230 disabled_filter_resolver,
231 disabled_tools,
232 disabled_skill_ids,
233 selected_skill_ids,
234 selected_skill_mode,
235 image_fallback,
236 gold_config,
237 guardian_config,
238 guardian_spawner,
239 bash_resume_hook,
240 bash_completion_sink,
241 app_data_dir,
242 run_budget,
243 } = params;
244
245 let mut builder = ExecuteRequestBuilder::new(initial_message, event_tx, cancel_token)
246 .model_roster(model_roster)
247 .gold_config(gold_config)
248 .guardian_config(guardian_config)
249 .guardian_spawner(guardian_spawner)
250 .bash_resume_hook(bash_resume_hook)
251 .bash_completion_sink(bash_completion_sink);
252
253 if let Some(run_budget) = run_budget {
254 builder = builder.run_budget(run_budget);
255 }
256
257 if let Some(tools) = tools {
258 builder = builder.tools(tools);
259 }
260 if let Some(provider_override) = provider_override {
261 builder = builder.provider_override(provider_override);
262 }
263 if let Some(reasoning_effort) = reasoning_effort {
264 builder = builder.reasoning_effort(reasoning_effort);
265 }
266 if let Some(disabled_filter_resolver) = disabled_filter_resolver {
267 builder = builder.disabled_filter_resolver(disabled_filter_resolver);
268 }
269 if let Some(auxiliary_model_resolver) = auxiliary_model_resolver {
270 builder = builder.auxiliary_model_resolver(auxiliary_model_resolver);
271 }
272 if let Some(disabled_tools) = disabled_tools {
273 builder = builder.disabled_tools(disabled_tools);
274 }
275 if let Some(disabled_skill_ids) = disabled_skill_ids {
276 builder = builder.disabled_skill_ids(disabled_skill_ids);
277 }
278 if let Some(selected_skill_ids) = selected_skill_ids {
279 builder = builder.selected_skill_ids(selected_skill_ids);
280 }
281 if let Some(selected_skill_mode) = selected_skill_mode {
282 builder = builder.selected_skill_mode(selected_skill_mode);
283 }
284 if let Some(image_fallback) = image_fallback {
285 builder = builder.image_fallback(image_fallback);
286 }
287 if let Some(app_data_dir) = app_data_dir {
288 builder = builder.app_data_dir(app_data_dir);
289 }
290
291 builder.build()
292}
293
294pub fn spawn_session_execution(args: SessionExecutionArgs) {
303 let span_session_id = args.session_id.clone();
304 let session_span = tracing::info_span!("agent_execution", session_id = %span_session_id);
305
306 tokio::spawn(
307 async move {
308 let SessionExecutionArgs {
309 agent,
310 session_id,
311 mut session,
312 tools_override,
313 provider_override,
314 model_roster,
315 reasoning_effort,
316 reasoning_effort_source,
317 auxiliary_model_resolver,
318 disabled_filter_resolver,
319 disabled_tools,
320 disabled_skill_ids,
321 selected_skill_ids,
322 selected_skill_mode,
323 cancel_token,
324 mpsc_tx,
325 image_fallback,
326 gold_config,
327 guardian_config,
328 guardian_spawner,
329 bash_resume_hook,
330 bash_completion_sink,
331 app_data_dir,
332 run_budget,
333 runners,
334 sessions_cache,
335 on_complete,
336 child_completion_handler,
337 } = args;
338
339 let model = model_roster.model.clone().unwrap_or_default();
343
344 let initial_message = initial_user_message_for_session(&session);
345 let selected_skill_ids =
346 selected_skill_ids.or_else(|| selected_skill_ids_for_session(&session));
347 let selected_skill_mode =
348 selected_skill_mode.or_else(|| selected_skill_mode_for_session(&session));
349
350 tracing::info!(
351 "[{}] Using resolved session model: {}, reasoning_effort={}, reasoning_source={}",
352 session_id,
353 model,
354 reasoning_effort
355 .map(ReasoningEffort::as_str)
356 .unwrap_or("none"),
357 reasoning_effort_source
358 );
359
360 crate::session_app::execution_prep::prepare_session_for_execution(
367 &mut session,
368 None,
369 Some(&model),
370 );
371
372 let system_prompt = system_prompt_for_session(&session);
373 if let Some(prompt) = system_prompt.as_ref() {
374 log_base_system_prompt_snapshot(&session_id, prompt);
375 }
376
377 let execute_request = build_execute_request(
378 initial_message,
379 mpsc_tx.clone(),
380 cancel_token,
381 ExecuteRequestParams {
382 tools: tools_override,
383 provider_override,
384 model_roster,
385 reasoning_effort,
386 auxiliary_model_resolver,
387 disabled_filter_resolver,
388 disabled_tools,
389 disabled_skill_ids,
390 selected_skill_ids,
391 selected_skill_mode,
392 image_fallback,
393 gold_config,
394 guardian_config,
395 guardian_spawner,
396 bash_resume_hook,
397 bash_completion_sink,
398 app_data_dir,
399 run_budget,
400 },
401 );
402
403 let result = {
410 use futures::FutureExt;
411 match std::panic::AssertUnwindSafe(agent.execute(&mut session, execute_request))
412 .catch_unwind()
413 .await
414 {
415 Ok(result) => result,
416 Err(panic) => {
417 let message = panic
418 .downcast_ref::<&str>()
419 .map(|s| (*s).to_string())
420 .or_else(|| panic.downcast_ref::<String>().cloned())
421 .unwrap_or_else(|| "non-string panic payload".to_string());
422 tracing::error!(
423 "[{}] agent execution panicked; finalizing as terminal error: {}",
424 session_id,
425 message
426 );
427 Err(AgentError::LLM(format!(
428 "agent execution panicked: {message}"
429 )))
430 }
431 }
432 };
433
434 if let Some(error_event) = terminal_error_event_for_result(&result) {
436 let _ = mpsc_tx.send(error_event).await;
437 }
438
439 let suspended_non_terminal = result.is_ok()
456 && session
457 .metadata
458 .get("runtime.suspend_reason")
459 .is_some_and(|reason| !reason.trim().is_empty());
460 match &result {
461 Ok(()) if suspended_non_terminal => {
462 session.set_last_run_status("suspended");
463 session.clear_last_run_error();
464 }
465 Ok(()) => {
466 session.set_last_run_status("completed");
467 session.clear_last_run_error();
468 }
469 Err(error) if error.is_cancelled() => {
470 session.set_last_run_status("cancelled");
471 session.set_last_run_error(error.to_string());
472 }
473 Err(error) => {
474 session.set_last_run_status("error");
475 session.set_last_run_error(error.to_string());
476 }
477 }
478
479 if let Some(on_complete) = on_complete {
483 on_complete(SessionExecutionOutcome::from_result(&result), &mut session).await;
484 }
485
486 if let Err(error) = agent.persistence().save_runtime_session(&mut session).await {
490 tracing::warn!("[{}] Failed to save session: {}", session_id, error);
491 }
492
493 finalize_runner(&runners, &session_id, &result).await;
499
500 let child_completion = child_completion_handler.filter(|_| {
509 session.kind == bamboo_agent_core::SessionKind::Child
510 && session.parent_session_id.is_some()
511 });
512 let parent_session_id = session.parent_session_id.clone();
513 let child_status = session.last_run_status();
514 let child_error = session.last_run_error();
515
516 sessions_cache.insert(
518 session_id.clone(),
519 Arc::new(parking_lot::RwLock::new(session)),
520 );
521
522 if let (Some(handler), Some(parent_session_id), Some(status)) =
523 (child_completion, parent_session_id, child_status)
524 {
525 use futures::FutureExt;
526 let completion = ChildCompletion {
527 parent_session_id: parent_session_id.clone(),
528 child_session_id: session_id.clone(),
529 status,
530 error: child_error,
531 completed_at: chrono::Utc::now(),
532 };
533 if std::panic::AssertUnwindSafe(handler.on_child_completed(completion))
534 .catch_unwind()
535 .await
536 .is_err()
537 {
538 tracing::error!(
539 %parent_session_id,
540 child_session_id = %session_id,
541 "child completion handler panicked on resumed-child terminal"
542 );
543 }
544 }
545
546 tracing::info!("[{}] Agent execution completed", session_id);
547 }
548 .instrument(session_span),
549 );
550}
551
552pub fn log_base_system_prompt_snapshot(session_id: &str, prompt: &str) {
554 tracing::info!(
555 "[{}] Base system prompt snapshot: len={} chars, has_skill={}, has_tool_guide={}, has_external_memory={}, has_task_list={}",
556 session_id,
557 prompt.len(),
558 prompt.contains(SKILL_CONTEXT_START_MARKER),
559 prompt.contains(TOOL_GUIDE_START_MARKER),
560 prompt.contains(EXTERNAL_MEMORY_START_MARKER),
561 prompt.contains(TASK_LIST_START_MARKER),
562 );
563
564 tracing::debug!(
565 "[{}] ========== BASE SYSTEM PROMPT SNAPSHOT ==========",
566 session_id
567 );
568 tracing::debug!("[{}] Snapshot length: {} chars", session_id, prompt.len());
569 tracing::debug!("[{}] -----------------------------------", session_id);
570 tracing::debug!("[{}] {}", session_id, prompt);
571 tracing::debug!(
572 "[{}] ========== END BASE SYSTEM PROMPT SNAPSHOT ==========",
573 session_id
574 );
575}
576
577pub fn terminal_error_event_for_result(result: &Result<(), AgentError>) -> Option<AgentEvent> {
579 match result {
580 Ok(_) => None,
581 Err(error) if error.is_cancelled() => Some(AgentEvent::Error {
582 message: "Agent execution cancelled by user".to_string(),
583 }),
584 Err(error) => Some(AgentEvent::Error {
585 message: error.to_string(),
586 }),
587 }
588}
589
590fn system_prompt_for_session(session: &Session) -> Option<String> {
593 session
594 .messages
595 .iter()
596 .find(|message| matches!(message.role, bamboo_agent_core::Role::System))
597 .map(|message| message.content.clone())
598}
599
600fn initial_user_message_for_session(session: &Session) -> String {
601 session
602 .messages
603 .last()
604 .filter(|message| matches!(message.role, bamboo_agent_core::Role::User))
605 .map(|message| message.content.clone())
606 .unwrap_or_default()
607}
608
609fn selected_skill_ids_for_session(session: &Session) -> Option<Vec<String>> {
610 session
611 .metadata
612 .get("selected_skill_ids")
613 .and_then(|raw| bamboo_skills::selection::parse_selected_skill_ids_metadata(raw))
614}
615
616fn selected_skill_mode_for_session(session: &Session) -> Option<String> {
617 let value = session
618 .metadata
619 .get("skill_mode")
620 .or_else(|| session.metadata.get("mode"))?;
621 let trimmed = value.trim();
622 if trimmed.is_empty() {
623 None
624 } else {
625 Some(trimmed.to_string())
626 }
627}