1use std::collections::{BTreeSet, HashMap};
7use std::future::Future;
8use std::pin::Pin;
9use std::sync::Arc;
10
11use tokio::sync::{broadcast, 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, DisabledFilterResolver, GoldConfig,
22 GuardianConfig, GuardianSpawner, ImageFallbackConfig,
23};
24use crate::runtime::execution::child_completion::ChildCompletion;
25use crate::runtime::execution::runner_lifecycle::{
26 finalize_rejected_runner_if_distinct, finalize_runner, finalize_runner_exact,
27 reserve_runner_core, ReserveOutcome, RunnerReservation,
28};
29use crate::runtime::execution::runner_state::AgentRunner;
30use crate::runtime::model_roster::ModelRoster;
31use crate::runtime::Agent;
32use crate::runtime::{ExecuteRequest, ExecuteRequestBuilder};
33use crate::session_activation::{
34 SessionActivationRouter, SessionRunRegistration, SessionRunRegistrationError,
35};
36
37pub type SessionCache = std::sync::Arc<
45 dashmap::DashMap<String, std::sync::Arc<parking_lot::RwLock<bamboo_agent_core::Session>>>,
46>;
47
48enum SessionExecutionActivationOwnership {
49 Unrouted,
51 UnpublishedActivation(Arc<SessionActivationRouter>),
54 RegistrationPending(Arc<SessionActivationRouter>),
59 Registered(SessionRunRegistration),
62}
63
64pub struct SessionExecutionReservation {
70 session_id: String,
71 run_id: String,
72 cancel_token: CancellationToken,
73 runners: Arc<RwLock<HashMap<String, AgentRunner>>>,
74 activation: SessionExecutionActivationOwnership,
75 armed: bool,
76}
77
78impl SessionExecutionReservation {
79 pub fn session_id(&self) -> &str {
80 &self.session_id
81 }
82
83 pub fn run_id(&self) -> &str {
84 &self.run_id
85 }
86
87 pub fn cancel_token(&self) -> &CancellationToken {
88 &self.cancel_token
89 }
90
91 pub(crate) fn from_activation_placeholder(
97 session_id: impl Into<String>,
98 reservation: RunnerReservation,
99 router: Arc<SessionActivationRouter>,
100 runners: Arc<RwLock<HashMap<String, AgentRunner>>>,
101 ) -> Self {
102 Self {
103 session_id: session_id.into(),
104 run_id: reservation.run_id,
105 cancel_token: reservation.cancel_token,
106 runners,
107 activation: SessionExecutionActivationOwnership::UnpublishedActivation(router),
108 armed: true,
109 }
110 }
111
112 pub(crate) fn from_pending_registration(
116 session_id: impl Into<String>,
117 reservation: RunnerReservation,
118 router: Option<Arc<SessionActivationRouter>>,
119 runners: Arc<RwLock<HashMap<String, AgentRunner>>>,
120 ) -> Self {
121 let activation = match router {
122 Some(router) => SessionExecutionActivationOwnership::RegistrationPending(router),
123 None => SessionExecutionActivationOwnership::Unrouted,
124 };
125 Self {
126 session_id: session_id.into(),
127 run_id: reservation.run_id,
128 cancel_token: reservation.cancel_token,
129 runners,
130 activation,
131 armed: true,
132 }
133 }
134
135 pub(crate) fn mark_activation_published(&mut self) {
138 let activation = std::mem::replace(
139 &mut self.activation,
140 SessionExecutionActivationOwnership::Unrouted,
141 );
142 self.activation = match activation {
143 SessionExecutionActivationOwnership::UnpublishedActivation(router) => {
144 SessionExecutionActivationOwnership::RegistrationPending(router)
145 }
146 other => other,
147 };
148 }
149
150 pub(crate) async fn rollback_unpublished_activation(mut self) {
155 self.armed = false;
156 self.cancel_token.cancel();
157 let activation = std::mem::replace(
158 &mut self.activation,
159 SessionExecutionActivationOwnership::Unrouted,
160 );
161 debug_assert!(matches!(
162 activation,
163 SessionExecutionActivationOwnership::UnpublishedActivation(_)
164 ));
165 remove_runner_exact(&self.runners, &self.session_id, &self.run_id).await;
166 }
167
168 pub async fn ensure_registered(&mut self) -> Result<(), SessionRunRegistrationError> {
174 let router = match &self.activation {
175 SessionExecutionActivationOwnership::RegistrationPending(router) => router.clone(),
176 SessionExecutionActivationOwnership::UnpublishedActivation(_) => {
177 unreachable!("unpublished activation entered an execution adapter")
178 }
179 SessionExecutionActivationOwnership::Unrouted
180 | SessionExecutionActivationOwnership::Registered(_) => return Ok(()),
181 };
182 match register_reserved_activation(router, &self.runners, &self.session_id, &self.run_id)
183 .await
184 {
185 Ok(registration) => {
186 self.activation = SessionExecutionActivationOwnership::Registered(registration);
187 Ok(())
188 }
189 Err(error) => {
190 self.disarm_after_registration_rejection(error.existing_run_id());
195 Err(error)
196 }
197 }
198 }
199
200 pub async fn abandon(mut self) {
204 self.armed = false;
205 self.cancel_token.cancel();
206 let activation = std::mem::replace(
207 &mut self.activation,
208 SessionExecutionActivationOwnership::Unrouted,
209 );
210 let runners = self.runners.clone();
211 let session_id = self.session_id.clone();
212 let run_id = self.run_id.clone();
213 let cleanup_session_id = session_id.clone();
214 let cleanup_run_id = run_id.clone();
215 let cleanup = tokio::spawn(async move {
216 cleanup_execution_reservation(activation, runners, cleanup_session_id, cleanup_run_id)
217 .await;
218 });
219 if let Err(error) = cleanup.await {
220 tracing::error!(
221 %session_id,
222 %run_id,
223 %error,
224 "detached execution-reservation cleanup failed"
225 );
226 }
227 }
228
229 fn disarm_after_registration_rejection(&mut self, existing_run_id: &str) {
230 if existing_run_id != self.run_id {
231 self.cancel_token.cancel();
232 }
233 self.armed = false;
234 }
235
236 pub(crate) fn matches_execution_target(
237 &self,
238 session_id: &str,
239 domain_session_id: &str,
240 runners: &Arc<RwLock<HashMap<String, AgentRunner>>>,
241 ) -> bool {
242 self.session_id == session_id
243 && domain_session_id == session_id
244 && Arc::ptr_eq(&self.runners, runners)
245 }
246
247 pub(crate) fn disarm_for_execution(
248 &mut self,
249 ) -> (CancellationToken, Option<SessionRunRegistration>) {
250 self.armed = false;
251 let registration = match std::mem::replace(
252 &mut self.activation,
253 SessionExecutionActivationOwnership::Unrouted,
254 ) {
255 SessionExecutionActivationOwnership::Registered(registration) => Some(registration),
256 SessionExecutionActivationOwnership::Unrouted => None,
257 SessionExecutionActivationOwnership::UnpublishedActivation(_) => {
258 unreachable!("unpublished activation cannot transfer to execution")
259 }
260 SessionExecutionActivationOwnership::RegistrationPending(_) => {
261 unreachable!("ensure_registered adopts every pending router registration")
262 }
263 };
264 (self.cancel_token.clone(), registration)
265 }
266}
267
268impl Drop for SessionExecutionReservation {
269 fn drop(&mut self) {
270 if !self.armed {
271 return;
272 }
273 self.armed = false;
274 self.cancel_token.cancel();
275 let activation = std::mem::replace(
276 &mut self.activation,
277 SessionExecutionActivationOwnership::Unrouted,
278 );
279 let runners = self.runners.clone();
280 let session_id = self.session_id.clone();
281 let run_id = self.run_id.clone();
282 if let Ok(runtime) = tokio::runtime::Handle::try_current() {
283 runtime.spawn(async move {
284 cleanup_execution_reservation(activation, runners, session_id, run_id).await;
285 });
286 }
287 }
288}
289
290pub enum SessionExecutionReserveOutcome {
292 Reserved(SessionExecutionReservation),
293 AlreadyRunning { run_id: String },
294}
295
296async fn register_reserved_activation(
297 router: Arc<SessionActivationRouter>,
298 runners: &Arc<RwLock<HashMap<String, AgentRunner>>>,
299 session_id: &str,
300 run_id: &str,
301) -> Result<SessionRunRegistration, SessionRunRegistrationError> {
302 let mut registration = match router.register_run(session_id, run_id).await {
303 Ok(registration) => registration,
304 Err(error) => {
305 let collision = Err(AgentError::LLM(error.to_string()));
306 finalize_rejected_runner_if_distinct(
307 runners,
308 session_id,
309 error.existing_run_id(),
310 run_id,
311 &collision,
312 )
313 .await;
314 return Err(error);
315 }
316 };
317 let cleanup_runners = runners.clone();
318 let cleanup_session_id = session_id.to_string();
319 let cleanup_run_id = run_id.to_string();
320 registration.set_abort_cleanup(move || async move {
321 let abandoned = Err(AgentError::Cancelled);
322 finalize_runner_exact(
323 &cleanup_runners,
324 &cleanup_session_id,
325 &cleanup_run_id,
326 &abandoned,
327 )
328 .await;
329 });
330 Ok(registration)
331}
332
333async fn cleanup_execution_reservation(
334 activation: SessionExecutionActivationOwnership,
335 runners: Arc<RwLock<HashMap<String, AgentRunner>>>,
336 session_id: String,
337 run_id: String,
338) {
339 match activation {
340 SessionExecutionActivationOwnership::Registered(registration) => {
341 registration.abandon().await;
342 }
343 SessionExecutionActivationOwnership::RegistrationPending(router) => {
344 match register_reserved_activation(router, &runners, &session_id, &run_id).await {
345 Ok(registration) => registration.abandon().await,
346 Err(error) => {
347 tracing::debug!(
348 %session_id,
349 %run_id,
350 existing_run_id = %error.existing_run_id(),
351 "abandoned activation placeholder was already adopted or superseded"
352 );
353 }
354 }
355 }
356 SessionExecutionActivationOwnership::UnpublishedActivation(_) => {
357 remove_runner_exact(&runners, &session_id, &run_id).await;
358 }
359 SessionExecutionActivationOwnership::Unrouted => {
360 let abandoned = Err(AgentError::Cancelled);
361 finalize_runner_exact(&runners, &session_id, &run_id, &abandoned).await;
362 }
363 }
364}
365
366async fn remove_runner_exact(
367 runners: &Arc<RwLock<HashMap<String, AgentRunner>>>,
368 session_id: &str,
369 run_id: &str,
370) -> bool {
371 let mut runners = runners.write().await;
372 if runners
373 .get(session_id)
374 .is_some_and(|runner| runner.run_id == run_id)
375 {
376 runners.remove(session_id);
377 true
378 } else {
379 false
380 }
381}
382
383pub async fn reserve_session_execution(
391 agent: &Arc<Agent>,
392 runners: &Arc<RwLock<HashMap<String, AgentRunner>>>,
393 senders: &Arc<RwLock<HashMap<String, broadcast::Sender<AgentEvent>>>>,
394 session_id: &str,
395 event_sender: &broadcast::Sender<AgentEvent>,
396) -> SessionExecutionReserveOutcome {
397 let reservation = match reserve_runner_core(runners, senders, session_id, event_sender).await {
398 ReserveOutcome::AlreadyRunning(run_id) => {
399 return SessionExecutionReserveOutcome::AlreadyRunning { run_id };
400 }
401 ReserveOutcome::Reserved(reservation) => reservation,
402 };
403
404 let mut execution_reservation = SessionExecutionReservation {
410 session_id: session_id.to_string(),
411 run_id: reservation.run_id,
412 cancel_token: reservation.cancel_token,
413 runners: runners.clone(),
414 activation: SessionExecutionActivationOwnership::Unrouted,
415 armed: true,
416 };
417
418 if let Some(router) = agent.activation_router().cloned() {
419 execution_reservation.activation =
420 SessionExecutionActivationOwnership::RegistrationPending(router);
421 if let Err(error) = execution_reservation.ensure_registered().await {
422 tracing::warn!(
423 %session_id,
424 attempted_run_id = %execution_reservation.run_id(),
425 existing_run_id = %error.existing_run_id(),
426 %error,
427 "runner reservation collided with an existing logical-session owner"
428 );
429 return SessionExecutionReserveOutcome::AlreadyRunning {
430 run_id: error.existing_run_id().to_string(),
431 };
432 }
433 }
434 SessionExecutionReserveOutcome::Reserved(execution_reservation)
435}
436
437pub fn read_cached_session(cache: &SessionCache, id: &str) -> Option<bamboo_agent_core::Session> {
445 cache
446 .get(id)
447 .map(|e| e.value().clone())
448 .map(|a| a.read().clone())
449}
450
451const SKILL_CONTEXT_START_MARKER: &str = "<!-- BAMBOO_SKILL_CONTEXT_START -->";
452const TOOL_GUIDE_START_MARKER: &str = "<!-- BAMBOO_TOOL_GUIDE_START -->";
453const EXTERNAL_MEMORY_START_MARKER: &str = "<!-- BAMBOO_EXTERNAL_MEMORY_START -->";
454const TASK_LIST_START_MARKER: &str = "<!-- BAMBOO_TASK_LIST_START -->";
455
456pub struct SessionExecutionOutcome {
463 pub success: bool,
465 pub cancelled: bool,
467 pub error: Option<String>,
469}
470
471impl SessionExecutionOutcome {
472 fn from_result(result: &Result<(), AgentError>) -> Self {
473 match result {
474 Ok(()) => Self {
475 success: true,
476 cancelled: false,
477 error: None,
478 },
479 Err(error) => Self {
480 success: false,
481 cancelled: error.is_cancelled(),
482 error: Some(error.to_string()),
483 },
484 }
485 }
486}
487
488pub type SessionCompletionHook = Box<
498 dyn for<'a> FnOnce(
499 SessionExecutionOutcome,
500 &'a mut Session,
501 ) -> Pin<Box<dyn Future<Output = ()> + Send + 'a>>
502 + Send,
503>;
504
505pub struct SessionExecutionArgs {
511 pub agent: Arc<Agent>,
513 pub session_id: String,
514 pub session: Session,
515 pub execution_reservation: SessionExecutionReservation,
518
519 pub tools_override: Option<Arc<dyn ToolExecutor>>,
521 pub provider_override: Option<Arc<dyn LLMProvider>>,
522 pub model_roster: ModelRoster,
526 pub reasoning_effort: Option<ReasoningEffort>,
527 pub reasoning_effort_source: String,
528 pub auxiliary_model_resolver:
529 Option<Arc<dyn Fn() -> crate::runtime::config::AuxiliaryModelConfig + Send + Sync>>,
530 pub disabled_filter_resolver: Option<DisabledFilterResolver>,
534 pub disabled_tools: Option<BTreeSet<String>>,
535 pub disabled_skill_ids: Option<BTreeSet<String>>,
536 pub selected_skill_ids: Option<Vec<String>>,
537 pub selected_skill_mode: Option<String>,
538 pub mpsc_tx: mpsc::Sender<AgentEvent>,
539 pub image_fallback: Option<ImageFallbackConfig>,
540 pub gold_config: Option<GoldConfig>,
541 pub guardian_config: Option<GuardianConfig>,
543 pub guardian_spawner: Option<Arc<dyn GuardianSpawner>>,
546 pub bash_resume_hook: Option<Arc<dyn BashResumeHook>>,
548 pub bash_completion_sink: Option<Arc<dyn BashCompletionSink>>,
550 pub app_data_dir: Option<std::path::PathBuf>,
551 pub run_budget: Option<bamboo_config::RunBudgetConfig>,
554
555 pub runners: Arc<RwLock<HashMap<String, AgentRunner>>>,
557 pub sessions_cache: SessionCache,
558
559 pub on_complete: Option<SessionCompletionHook>,
562
563 pub child_completion_handler: Option<Arc<dyn super::ChildCompletionHandler>>,
575}
576
577struct ExecuteRequestParams {
587 tools: Option<Arc<dyn ToolExecutor>>,
588 provider_override: Option<Arc<dyn LLMProvider>>,
589 model_roster: ModelRoster,
590 reasoning_effort: Option<ReasoningEffort>,
591 auxiliary_model_resolver: Option<Arc<dyn Fn() -> AuxiliaryModelConfig + Send + Sync>>,
592 disabled_filter_resolver: Option<DisabledFilterResolver>,
593 disabled_tools: Option<BTreeSet<String>>,
594 disabled_skill_ids: Option<BTreeSet<String>>,
595 selected_skill_ids: Option<Vec<String>>,
596 selected_skill_mode: Option<String>,
597 image_fallback: Option<ImageFallbackConfig>,
598 gold_config: Option<GoldConfig>,
599 guardian_config: Option<GuardianConfig>,
600 guardian_spawner: Option<Arc<dyn GuardianSpawner>>,
601 bash_resume_hook: Option<Arc<dyn BashResumeHook>>,
602 bash_completion_sink: Option<Arc<dyn BashCompletionSink>>,
603 app_data_dir: Option<std::path::PathBuf>,
604 run_budget: Option<bamboo_config::RunBudgetConfig>,
605}
606
607fn build_execute_request(
614 initial_message: String,
615 event_tx: mpsc::Sender<AgentEvent>,
616 cancel_token: CancellationToken,
617 params: ExecuteRequestParams,
618) -> ExecuteRequest {
619 let ExecuteRequestParams {
620 tools,
621 provider_override,
622 model_roster,
623 reasoning_effort,
624 auxiliary_model_resolver,
625 disabled_filter_resolver,
626 disabled_tools,
627 disabled_skill_ids,
628 selected_skill_ids,
629 selected_skill_mode,
630 image_fallback,
631 gold_config,
632 guardian_config,
633 guardian_spawner,
634 bash_resume_hook,
635 bash_completion_sink,
636 app_data_dir,
637 run_budget,
638 } = params;
639
640 let mut builder = ExecuteRequestBuilder::new(initial_message, event_tx, cancel_token)
641 .model_roster(model_roster)
642 .gold_config(gold_config)
643 .guardian_config(guardian_config)
644 .guardian_spawner(guardian_spawner)
645 .bash_resume_hook(bash_resume_hook)
646 .bash_completion_sink(bash_completion_sink);
647
648 if let Some(run_budget) = run_budget {
649 builder = builder.run_budget(run_budget);
650 }
651
652 if let Some(tools) = tools {
653 builder = builder.tools(tools);
654 }
655 if let Some(provider_override) = provider_override {
656 builder = builder.provider_override(provider_override);
657 }
658 if let Some(reasoning_effort) = reasoning_effort {
659 builder = builder.reasoning_effort(reasoning_effort);
660 }
661 if let Some(disabled_filter_resolver) = disabled_filter_resolver {
662 builder = builder.disabled_filter_resolver(disabled_filter_resolver);
663 }
664 if let Some(auxiliary_model_resolver) = auxiliary_model_resolver {
665 builder = builder.auxiliary_model_resolver(auxiliary_model_resolver);
666 }
667 if let Some(disabled_tools) = disabled_tools {
668 builder = builder.disabled_tools(disabled_tools);
669 }
670 if let Some(disabled_skill_ids) = disabled_skill_ids {
671 builder = builder.disabled_skill_ids(disabled_skill_ids);
672 }
673 if let Some(selected_skill_ids) = selected_skill_ids {
674 builder = builder.selected_skill_ids(selected_skill_ids);
675 }
676 if let Some(selected_skill_mode) = selected_skill_mode {
677 builder = builder.selected_skill_mode(selected_skill_mode);
678 }
679 if let Some(image_fallback) = image_fallback {
680 builder = builder.image_fallback(image_fallback);
681 }
682 if let Some(app_data_dir) = app_data_dir {
683 builder = builder.app_data_dir(app_data_dir);
684 }
685
686 builder.build()
687}
688
689pub fn spawn_session_execution(args: SessionExecutionArgs) {
698 let span_session_id = args.session_id.clone();
699 let session_span = tracing::info_span!("agent_execution", session_id = %span_session_id);
700
701 tokio::spawn(
702 async move {
703 let SessionExecutionArgs {
704 agent,
705 session_id,
706 mut session,
707 mut execution_reservation,
708 tools_override,
709 provider_override,
710 model_roster,
711 reasoning_effort,
712 reasoning_effort_source,
713 auxiliary_model_resolver,
714 disabled_filter_resolver,
715 disabled_tools,
716 disabled_skill_ids,
717 selected_skill_ids,
718 selected_skill_mode,
719 mpsc_tx,
720 image_fallback,
721 gold_config,
722 guardian_config,
723 guardian_spawner,
724 bash_resume_hook,
725 bash_completion_sink,
726 app_data_dir,
727 run_budget,
728 runners,
729 sessions_cache,
730 on_complete,
731 child_completion_handler,
732 } = args;
733
734 if !execution_reservation.matches_execution_target(&session_id, &session.id, &runners) {
735 tracing::error!(
736 %session_id,
737 domain_session_id = %session.id,
738 reservation_session_id = %execution_reservation.session_id(),
739 run_id = %execution_reservation.run_id(),
740 same_runner_registry = Arc::ptr_eq(&execution_reservation.runners, &runners),
741 "refusing mismatched session execution reservation"
742 );
743 execution_reservation.abandon().await;
744 return;
745 }
746 if let Err(error) = execution_reservation.ensure_registered().await {
747 tracing::warn!(
748 %session_id,
749 run_id = %execution_reservation.run_id(),
750 %error,
751 "session execution could not adopt its router activation owner"
752 );
753 return;
754 }
755 let (cancel_token, mut activation_registration) =
756 execution_reservation.disarm_for_execution();
757
758 let model = model_roster.model.clone().unwrap_or_default();
762
763 let initial_message = initial_user_message_for_session(&session);
764 let selected_skill_ids =
765 selected_skill_ids.or_else(|| selected_skill_ids_for_session(&session));
766 let selected_skill_mode =
767 selected_skill_mode.or_else(|| selected_skill_mode_for_session(&session));
768
769 tracing::info!(
770 "[{}] Using resolved session model: {}, reasoning_effort={}, reasoning_source={}",
771 session_id,
772 model,
773 reasoning_effort
774 .map(ReasoningEffort::as_str)
775 .unwrap_or("none"),
776 reasoning_effort_source
777 );
778
779 crate::session_app::execution_prep::prepare_session_for_execution(
786 &mut session,
787 None,
788 Some(&model),
789 );
790
791 let system_prompt = system_prompt_for_session(&session);
792 if let Some(prompt) = system_prompt.as_ref() {
793 log_base_system_prompt_snapshot(&session_id, prompt);
794 }
795
796 let execute_request = build_execute_request(
797 initial_message,
798 mpsc_tx.clone(),
799 cancel_token,
800 ExecuteRequestParams {
801 tools: tools_override,
802 provider_override,
803 model_roster,
804 reasoning_effort,
805 auxiliary_model_resolver,
806 disabled_filter_resolver,
807 disabled_tools,
808 disabled_skill_ids,
809 selected_skill_ids,
810 selected_skill_mode,
811 image_fallback,
812 gold_config,
813 guardian_config,
814 guardian_spawner,
815 bash_resume_hook,
816 bash_completion_sink,
817 app_data_dir,
818 run_budget,
819 },
820 );
821
822 let result = {
829 use futures::FutureExt;
830 match std::panic::AssertUnwindSafe(agent.execute(&mut session, execute_request))
831 .catch_unwind()
832 .await
833 {
834 Ok(result) => result,
835 Err(panic) => {
836 let message = panic
837 .downcast_ref::<&str>()
838 .map(|s| (*s).to_string())
839 .or_else(|| panic.downcast_ref::<String>().cloned())
840 .unwrap_or_else(|| "non-string panic payload".to_string());
841 tracing::error!(
842 "[{}] agent execution panicked; finalizing as terminal error: {}",
843 session_id,
844 message
845 );
846 Err(AgentError::LLM(format!(
847 "agent execution panicked: {message}"
848 )))
849 }
850 }
851 };
852
853 if let Some(error_event) = terminal_error_event_for_result(&result) {
855 let _ = mpsc_tx.send(error_event).await;
856 }
857
858 let suspended_non_terminal = result.is_ok()
875 && session
876 .metadata
877 .get("runtime.suspend_reason")
878 .is_some_and(|reason| !reason.trim().is_empty());
879 match &result {
880 Ok(()) if suspended_non_terminal => {
881 session.set_last_run_status("suspended");
882 session.clear_last_run_error();
883 }
884 Ok(()) => {
885 session.set_last_run_status("completed");
886 session.clear_last_run_error();
887 }
888 Err(error) if error.is_cancelled() => {
889 session.set_last_run_status("cancelled");
890 session.set_last_run_error(error.to_string());
891 }
892 Err(error) => {
893 session.set_last_run_status("error");
894 session.set_last_run_error(error.to_string());
895 }
896 }
897
898 if let Some(on_complete) = on_complete {
902 on_complete(SessionExecutionOutcome::from_result(&result), &mut session).await;
903 }
904
905 let executed_admitted_generation = session
909 .session_inbox_admission()
910 .map_or(0, |state| state.last_admitted_sequence);
911 let legacy_migration =
912 crate::runtime::runner::state_bridge::migrate_legacy_pending_only(
913 &mut session,
914 Some(agent.storage()),
915 Some(agent.persistence()),
916 agent.session_inbox(),
917 )
918 .await;
919 let pending_boundary_generation = session
920 .session_inbox_admission()
921 .and_then(|state| state.pending_activation_generation());
922 let pending_generation = match (
923 pending_boundary_generation,
924 legacy_migration.highest_generation,
925 ) {
926 (Some(left), Some(right)) => Some(left.max(right)),
927 (left, right) => left.or(right),
928 };
929 if let (Some(generation), Some(router)) =
930 (pending_generation, agent.activation_router())
931 {
932 let activation_ready = if let Some(inbox) = agent.session_inbox() {
933 match inbox
934 .mark_activation_eligible(
935 &session_id,
936 generation,
937 bamboo_domain::SessionActivationPolicy::InterruptSpecificWait,
938 )
939 .await
940 {
941 Ok(()) => true,
942 Err(error) => {
943 tracing::error!(
944 session_id = %session_id,
945 %error,
946 "failed to persist unadmitted SessionInbox activation watermark"
947 );
948 false
949 }
950 }
951 } else {
952 false
953 };
954 if activation_ready {
955 if let Err(error) = bamboo_domain::SessionActivationPort::request_activation(
956 router.as_ref(),
957 &session_id,
958 generation,
959 )
960 .await
961 {
962 tracing::error!(
963 session_id = %session_id,
964 %error,
965 "failed to hand unadmitted SessionInbox generation to activation router"
966 );
967 }
968 }
969 }
970
971 if let Some(registration) = activation_registration.as_mut() {
975 registration.begin_finalization().await;
976 }
977
978 if let Err(error) = agent.persistence().save_runtime_session(&mut session).await {
982 tracing::warn!("[{}] Failed to save session: {}", session_id, error);
983 }
984
985 finalize_runner(&runners, &session_id, &result).await;
991
992 let finalization = if let Some(registration) = activation_registration.take() {
993 registration.finish(executed_admitted_generation).await
994 } else {
995 Ok(None)
996 };
997 if let Err(error) = finalization {
998 tracing::error!(
1002 session_id = %session_id,
1003 %error,
1004 "failed to activate successor for finalization-racing SessionInbox delivery"
1005 );
1006 }
1007
1008 let child_completion = child_completion_handler.filter(|_| {
1017 session.kind == bamboo_agent_core::SessionKind::Child
1018 && session.parent_session_id.is_some()
1019 });
1020 let parent_session_id = session.parent_session_id.clone();
1021 let child_status = session.last_run_status();
1022 let child_error = session.last_run_error();
1023
1024 sessions_cache.insert(
1026 session_id.clone(),
1027 Arc::new(parking_lot::RwLock::new(session)),
1028 );
1029
1030 if let (Some(handler), Some(parent_session_id), Some(status)) =
1031 (child_completion, parent_session_id, child_status)
1032 {
1033 use futures::FutureExt;
1034 let completion = ChildCompletion {
1035 parent_session_id: parent_session_id.clone(),
1036 child_session_id: session_id.clone(),
1037 status,
1038 error: child_error,
1039 completed_at: chrono::Utc::now(),
1040 };
1041 if std::panic::AssertUnwindSafe(handler.on_child_completed(completion))
1042 .catch_unwind()
1043 .await
1044 .is_err()
1045 {
1046 tracing::error!(
1047 %parent_session_id,
1048 child_session_id = %session_id,
1049 "child completion handler panicked on resumed-child terminal"
1050 );
1051 }
1052 }
1053
1054 tracing::info!("[{}] Agent execution completed", session_id);
1055 }
1056 .instrument(session_span),
1057 );
1058}
1059
1060pub fn log_base_system_prompt_snapshot(session_id: &str, prompt: &str) {
1062 tracing::info!(
1063 "[{}] Base system prompt snapshot: len={} chars, has_skill={}, has_tool_guide={}, has_external_memory={}, has_task_list={}",
1064 session_id,
1065 prompt.len(),
1066 prompt.contains(SKILL_CONTEXT_START_MARKER),
1067 prompt.contains(TOOL_GUIDE_START_MARKER),
1068 prompt.contains(EXTERNAL_MEMORY_START_MARKER),
1069 prompt.contains(TASK_LIST_START_MARKER),
1070 );
1071
1072 tracing::debug!(
1073 "[{}] ========== BASE SYSTEM PROMPT SNAPSHOT ==========",
1074 session_id
1075 );
1076 tracing::debug!("[{}] Snapshot length: {} chars", session_id, prompt.len());
1077 tracing::debug!("[{}] -----------------------------------", session_id);
1078 tracing::debug!("[{}] {}", session_id, prompt);
1079 tracing::debug!(
1080 "[{}] ========== END BASE SYSTEM PROMPT SNAPSHOT ==========",
1081 session_id
1082 );
1083}
1084
1085pub fn terminal_error_event_for_result(result: &Result<(), AgentError>) -> Option<AgentEvent> {
1087 match result {
1088 Ok(_) => None,
1089 Err(error) if error.is_cancelled() => Some(AgentEvent::Error {
1090 message: "Agent execution cancelled by user".to_string(),
1091 }),
1092 Err(error) => Some(AgentEvent::Error {
1093 message: error.to_string(),
1094 }),
1095 }
1096}
1097
1098fn system_prompt_for_session(session: &Session) -> Option<String> {
1101 session
1102 .messages
1103 .iter()
1104 .find(|message| matches!(message.role, bamboo_agent_core::Role::System))
1105 .map(|message| message.content.clone())
1106}
1107
1108fn initial_user_message_for_session(session: &Session) -> String {
1109 session
1110 .messages
1111 .last()
1112 .filter(|message| matches!(message.role, bamboo_agent_core::Role::User))
1113 .map(|message| message.content.clone())
1114 .unwrap_or_default()
1115}
1116
1117fn selected_skill_ids_for_session(session: &Session) -> Option<Vec<String>> {
1118 session
1119 .metadata
1120 .get("selected_skill_ids")
1121 .and_then(|raw| bamboo_skills::selection::parse_selected_skill_ids_metadata(raw))
1122}
1123
1124fn selected_skill_mode_for_session(session: &Session) -> Option<String> {
1125 let value = session
1126 .metadata
1127 .get("skill_mode")
1128 .or_else(|| session.metadata.get("mode"))?;
1129 let trimmed = value.trim();
1130 if trimmed.is_empty() {
1131 None
1132 } else {
1133 Some(trimmed.to_string())
1134 }
1135}
1136
1137#[cfg(test)]
1138mod reservation_tests {
1139 use super::*;
1140 use crate::runtime::execution::runner_state::AgentStatus;
1141
1142 #[test]
1143 fn reservation_target_requires_domain_id_and_exact_runner_registry() {
1144 let runners = Arc::new(RwLock::new(HashMap::new()));
1145 let other_runners = Arc::new(RwLock::new(HashMap::new()));
1146 let mut reservation = SessionExecutionReservation {
1147 session_id: "session-a".to_string(),
1148 run_id: "run-a".to_string(),
1149 cancel_token: CancellationToken::new(),
1150 runners: runners.clone(),
1151 activation: SessionExecutionActivationOwnership::Unrouted,
1152 armed: true,
1153 };
1154
1155 assert!(reservation.matches_execution_target("session-a", "session-a", &runners));
1156 assert!(!reservation.matches_execution_target("session-b", "session-a", &runners));
1157 assert!(!reservation.matches_execution_target("session-a", "session-b", &runners));
1158 assert!(!reservation.matches_execution_target("session-a", "session-a", &other_runners));
1159 reservation.armed = false;
1160 }
1161
1162 #[tokio::test]
1163 async fn forced_same_run_registration_rejection_never_cancels_live_owner() {
1164 let runners = Arc::new(RwLock::new(HashMap::new()));
1165 let mut runner = AgentRunner::new();
1166 runner.status = AgentStatus::Running;
1167 let run_id = runner.run_id.clone();
1168 let live_cancel_token = runner.cancel_token.clone();
1169 runners.write().await.insert("same-run".to_string(), runner);
1170
1171 let router = SessionActivationRouter::new();
1172 let mut live_registration = router
1173 .register_run("same-run", &run_id)
1174 .await
1175 .expect("first registration owns the run");
1176 let mut rejected = SessionExecutionReservation {
1177 session_id: "same-run".to_string(),
1178 run_id: run_id.clone(),
1179 cancel_token: live_cancel_token.clone(),
1180 runners: runners.clone(),
1181 activation: SessionExecutionActivationOwnership::RegistrationPending(router.clone()),
1182 armed: true,
1183 };
1184
1185 let error = match rejected.ensure_registered().await {
1186 Ok(()) => panic!("a duplicate registration for the same run must be rejected"),
1187 Err(error) => error,
1188 };
1189 assert_eq!(error.existing_run_id(), run_id);
1190 drop(rejected);
1191
1192 assert!(!live_cancel_token.is_cancelled());
1193 assert!(matches!(
1194 runners
1195 .read()
1196 .await
1197 .get("same-run")
1198 .map(|runner| &runner.status),
1199 Some(AgentStatus::Running)
1200 ));
1201 assert!(router.owns_run("same-run", &run_id).await);
1202 live_registration.begin_finalization().await;
1203 assert_eq!(live_registration.finish(0).await.unwrap(), None);
1204 }
1205}