use std::sync::Arc;
use tokio::sync::mpsc;
use tokio_util::sync::CancellationToken;
use tracing::Instrument;
use crate::runtime::config::AgentLoopConfig;
use bamboo_agent_core::tools::ToolExecutor;
use bamboo_agent_core::{AgentEvent, Session};
use bamboo_llm::LLMProvider;
mod gold;
mod pipeline;
mod startup;
use pipeline::run_pipeline;
use startup::{initialize_loop_state, LoopRunState};
pub(crate) async fn run_agent_loop_with_config(
session: &mut Session,
initial_message: String,
event_tx: mpsc::Sender<AgentEvent>,
llm: Arc<dyn LLMProvider>,
tools: Arc<dyn ToolExecutor>,
cancel_token: CancellationToken,
config: AgentLoopConfig,
) -> super::Result<()> {
let session_span = tracing::info_span!("agent_loop", session_id = %session.id);
async {
let mut state: LoopRunState = initialize_loop_state(
session,
initial_message.as_str(),
&config,
tools.as_ref(),
&event_tx,
)
.await?;
let pipeline_result = run_pipeline(
session,
&event_tx,
llm,
tools,
&cancel_token,
&config,
&mut state,
)
.await;
let sent_complete = match pipeline_result {
Ok(sent_complete) => sent_complete,
Err(error) => {
if let Some(skill_manager) = config.skill_manager.as_ref() {
let workspace = session.workspace_path_meta().map(std::path::PathBuf::from);
if let Err(release_error) = skill_manager
.release_activation_for_workspace(&state.session_id, workspace.as_deref())
.await
{
tracing::warn!(
"[{}] Failed to release errored workflow activation snapshot: {}",
state.session_id,
release_error
);
}
}
return Err(error);
}
};
super::session_finalize::finalize_session(
state.task_context,
session,
&event_tx,
&state.session_id,
&config,
state.metrics_collector.as_ref(),
sent_complete,
&mut state.runtime_state,
)
.await;
Ok(())
}
.instrument(session_span)
.await
}