use super::*;
impl AgentSession {
pub fn has_queue(&self) -> bool {
QueueControl::from_session(self).has_queue()
}
pub async fn set_lane_handler(
&self,
lane: SessionLane,
config: LaneHandlerConfig,
) -> crate::error::Result<()> {
let _mutation = self.close_handle.extension_mutation.lock().await;
if self.is_closed() {
return Err(crate::error::CodeError::SessionClosed {
session_id: self.session_id.clone(),
});
}
QueueControl::from_session(self)
.set_lane_handler(lane, config)
.await;
if self.is_closed() {
return Err(crate::error::CodeError::SessionClosed {
session_id: self.session_id.clone(),
});
}
Ok(())
}
pub async fn complete_external_task(&self, task_id: &str, result: ExternalTaskResult) -> bool {
QueueControl::from_session(self)
.complete_external_task(task_id, result)
.await
}
pub async fn pending_external_tasks(&self) -> Vec<ExternalTask> {
QueueControl::from_session(self)
.pending_external_tasks()
.await
}
pub async fn queue_stats(&self) -> SessionQueueStats {
QueueControl::from_session(self).stats().await
}
pub async fn queue_metrics(&self) -> Option<MetricsSnapshot> {
QueueControl::from_session(self).metrics().await
}
pub async fn dead_letters(&self) -> Vec<DeadLetter> {
QueueControl::from_session(self).dead_letters().await
}
pub fn register_agent_dir(&self, dir: &std::path::Path) -> crate::error::Result<usize> {
let agents = crate::subagent::load_agents_from_dir(dir);
self.close_handle.mutate_immediate(|| {
for agent in &agents {
self.ensure_compatibility_name_available(
crate::capability::CapabilityKind::Agent,
&agent.name,
)?;
}
let count = agents.len();
for agent in agents {
tracing::info!(
session_id = %self.session_id,
agent = agent.name,
dir = %dir.display(),
"Dynamically registered agent"
);
self.agent_registry.register(agent);
}
Ok(count)
})?
}
pub fn register_worker_agent(
&self,
spec: crate::subagent::WorkerAgentSpec,
) -> crate::error::Result<crate::subagent::AgentDefinition> {
self.close_handle.mutate_immediate(|| {
self.ensure_compatibility_name_available(
crate::capability::CapabilityKind::Agent,
&spec.name,
)?;
Ok(SessionExtensionRuntime::from_session(self).register_worker_agent(spec))
})?
}
pub fn register_worker_agents<I>(
&self,
specs: I,
) -> crate::error::Result<Vec<crate::subagent::AgentDefinition>>
where
I: IntoIterator<Item = crate::subagent::WorkerAgentSpec>,
{
let specs = specs.into_iter().collect::<Vec<_>>();
self.close_handle.mutate_immediate(|| {
for spec in &specs {
self.ensure_compatibility_name_available(
crate::capability::CapabilityKind::Agent,
&spec.name,
)?;
}
Ok(SessionExtensionRuntime::from_session(self).register_worker_agents(specs))
})?
}
pub fn add_skill(&self, skill: Arc<crate::skills::Skill>) -> crate::error::Result<()> {
self.close_handle.mutate_immediate(|| {
self.ensure_compatibility_name_available(
crate::capability::CapabilityKind::Skill,
&skill.name,
)?;
SessionExtensionRuntime::from_session(self).add_skill(skill)
})?
}
pub fn remove_skill(&self, name: &str) -> crate::error::Result<()> {
self.close_handle
.mutate_immediate(|| SessionExtensionRuntime::from_session(self).remove_skill(name))
}
pub fn skill_names(&self) -> Vec<String> {
self.close_handle.skill_registry.list()
}
pub async fn add_mcp_server(
&self,
config: crate::mcp::McpServerConfig,
) -> crate::error::Result<usize> {
SessionExtensionRuntime::from_session(self)
.add_mcp_server(config)
.await
}
#[cfg(feature = "serve")]
pub(crate) fn tool_executor(&self) -> &Arc<crate::tools::ToolExecutor> {
&self.tool_executor
}
pub fn register_dynamic_tool(
&self,
tool: Arc<dyn crate::tools::Tool>,
) -> crate::error::Result<()> {
self.close_handle.mutate_immediate(|| {
self.ensure_compatibility_name_available(
crate::capability::CapabilityKind::Tool,
tool.name(),
)?;
self.tool_executor.register_dynamic_tool(tool);
Ok(())
})?
}
pub fn register_dynamic_workflow_runtime(&self) -> crate::error::Result<()> {
self.close_handle.mutate_immediate(|| {
self.ensure_compatibility_name_available(
crate::capability::CapabilityKind::Tool,
"dynamic_workflow",
)?;
crate::tools::register_dynamic_workflow(self.tool_executor.registry());
Ok(())
})?
}
pub fn unregister_dynamic_tool(&self, name: &str) -> crate::error::Result<()> {
self.close_handle
.mutate_immediate(|| self.tool_executor.unregister_dynamic_tool(name))
}
pub async fn remove_mcp_server(&self, server_name: &str) -> crate::error::Result<()> {
SessionExtensionRuntime::from_session(self)
.remove_mcp_server(server_name)
.await
}
pub async fn mcp_status(
&self,
) -> std::collections::HashMap<String, crate::mcp::McpServerStatus> {
SessionExtensionRuntime::from_session(self)
.mcp_status()
.await
}
pub fn capability_catalog_stamp(&self) -> crate::capability::CapabilityCatalogStamp {
self.capability_catalog.current_stamp()
}
pub fn ensure_recovery_capability_binding(
&self,
expected: &crate::capability::RunCapabilityBindingV1,
) -> std::result::Result<(), crate::capability::RunCapabilityBindingError> {
super::agent_loop_runtime::validate_run_capability_binding(self, expected)
}
pub async fn apply_capability_batch(
&self,
batch: crate::capability::SessionCapabilityBatch,
cancellation: tokio_util::sync::CancellationToken,
) -> std::result::Result<
crate::capability::CapabilityCommitReceipt,
crate::capability::CapabilityRuntimeError,
> {
let _mutation = self.close_handle.extension_mutation.lock().await;
if self.is_closed() {
return Err(crate::capability::CapabilityRuntimeError::SessionClosed);
}
let preparation_cancellation = tokio_util::sync::CancellationToken::new();
let prepared = tokio::select! {
biased;
_ = self.session_cancel.cancelled() => {
preparation_cancellation.cancel();
return Err(if self.is_closed() {
crate::capability::CapabilityRuntimeError::SessionClosed
} else {
crate::capability::CapabilityRuntimeError::Cancelled
});
}
_ = cancellation.cancelled() => {
preparation_cancellation.cancel();
return Err(crate::capability::CapabilityRuntimeError::Cancelled);
}
result = batch.prepare(
&self.capability_catalog,
preparation_cancellation.clone(),
) => result?,
};
self.ensure_projected_mcp_server_names_available(prepared.projection()?)
.await?;
let _publication = self
.close_handle
.immediate_extension_mutation
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
if self.is_closed() {
return Err(crate::capability::CapabilityRuntimeError::SessionClosed);
}
if cancellation.is_cancelled() || self.session_cancel.is_cancelled() {
return Err(crate::capability::CapabilityRuntimeError::Cancelled);
}
let command_registry = self
.command_registry
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
let current_projection = self.capability_catalog.pin();
super::agent_loop_runtime::validate_capability_projection_runtime(
self,
prepared.projection()?,
&command_registry,
)?;
super::agent_loop_runtime::validate_capability_projection_transition(
self,
current_projection.projection(),
prepared.projection()?,
)?;
prepared.commit()
}
pub async fn bootstrap_recovery_capability_batch(
&self,
expected: &crate::capability::RunCapabilityBindingV1,
batch: crate::capability::SessionCapabilityBatch,
cancellation: tokio_util::sync::CancellationToken,
) -> std::result::Result<
crate::capability::CapabilityCommitReceipt,
crate::capability::CapabilityRuntimeError,
> {
let _mutation = self.close_handle.extension_mutation.lock().await;
if self.is_closed() {
return Err(crate::capability::CapabilityRuntimeError::SessionClosed);
}
let target_ceiling = self.capability_run_ceiling(batch.target())?;
expected
.ensure_matches(batch.target(), &target_ceiling)
.map_err(
|error| crate::capability::CapabilityRuntimeError::RecoveryBinding {
message: error.to_string(),
},
)?;
let preparation_cancellation = tokio_util::sync::CancellationToken::new();
let prepared = tokio::select! {
biased;
_ = self.session_cancel.cancelled() => {
preparation_cancellation.cancel();
return Err(if self.is_closed() {
crate::capability::CapabilityRuntimeError::SessionClosed
} else {
crate::capability::CapabilityRuntimeError::Cancelled
});
}
_ = cancellation.cancelled() => {
preparation_cancellation.cancel();
return Err(crate::capability::CapabilityRuntimeError::Cancelled);
}
result = batch.prepare_recovery_bootstrap(
&self.capability_catalog,
preparation_cancellation.clone(),
) => result?,
};
self.ensure_projected_mcp_server_names_available(prepared.projection()?)
.await?;
let _publication = self
.close_handle
.immediate_extension_mutation
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
if self.is_closed() {
return Err(crate::capability::CapabilityRuntimeError::SessionClosed);
}
if cancellation.is_cancelled() || self.session_cancel.is_cancelled() {
return Err(crate::capability::CapabilityRuntimeError::Cancelled);
}
let command_registry = self
.command_registry
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
let current_projection = self.capability_catalog.pin();
super::agent_loop_runtime::validate_capability_projection_runtime(
self,
prepared.projection()?,
&command_registry,
)?;
super::agent_loop_runtime::validate_capability_projection_transition(
self,
current_projection.projection(),
prepared.projection()?,
)?;
prepared.commit()
}
pub async fn drain_capability_cleanup(&self) -> crate::capability::CapabilityCleanupReport {
self.capability_catalog.drain_cleanup().await
}
#[cfg(test)]
pub(crate) async fn admit_capability_run(
&self,
) -> std::result::Result<
crate::capability::SessionCapabilityRun,
crate::capability::CapabilityRuntimeError,
> {
let projection = {
let _admission = self
.close_handle
.immediate_extension_mutation
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
if self.is_closed() {
return Err(crate::capability::CapabilityRuntimeError::SessionClosed);
}
self.capability_catalog.pin()
};
let ceiling = self.capability_run_ceiling(projection.projection().set())?;
crate::capability::SessionCapabilityRun::admit(
projection,
"active",
"active",
ceiling,
self.session_cancel.child_token(),
)
.await
}
pub(super) fn capability_run_ceiling(
&self,
set: &crate::capability::CapabilitySet,
) -> std::result::Result<
crate::capability::CapabilityCeiling,
crate::capability::CapabilityRuntimeError,
> {
let mut governance = crate::capability::GovernanceCapabilityCeiling::none_required();
if self.config.permission_checker.is_some() || self.config.permission_policy.is_some() {
governance = governance.require_permission_guard();
}
if self.config.confirmation_manager.is_some() || self.config.confirmation_policy.is_some() {
governance = governance.require_confirmation_guard();
}
if self.config.security_provider.is_some() {
governance = governance.require_security_guard();
}
if self.config.budget_guard.is_some() || self.budget_guard().is_some() {
governance = governance.require_budget_guard();
}
if self.config.enforce_active_skill_tool_restrictions {
governance = governance.require_active_skill_restrictions();
}
let execution = crate::capability::CapabilityExecutionCeiling::new(
self.config.max_tool_rounds,
self.config.max_parallel_tasks,
self.config.tool_timeout_ms,
self.config.llm_api_timeout_ms,
self.config.max_execution_time_ms,
)?;
crate::capability::CapabilityCeiling::all(
set,
crate::capability::WorkspaceCapabilityCeiling::all(),
governance,
execution,
)
.map_err(Into::into)
}
pub(super) fn ensure_compatibility_name_available(
&self,
kind: crate::capability::CapabilityKind,
public_name: &str,
) -> crate::error::Result<()> {
let projection = self.capability_catalog.pin();
if projection
.projection()
.iter()
.any(|(_, value)| match value {
crate::capability::CapabilityValue::Mcp(binding)
if kind == crate::capability::CapabilityKind::Tool =>
{
binding.contains_public_tool_name(public_name)
}
crate::capability::CapabilityValue::Agent(agent) => {
kind == crate::capability::CapabilityKind::Agent
&& crate::subagent::agent_names_conflict(&agent.name, public_name)
}
_ => value.kind() == kind && value.public_name() == Some(public_name),
})
{
return Err(
crate::capability::CapabilityRuntimeError::RuntimeNameConflict {
kind,
public_name: public_name.to_owned(),
}
.into(),
);
}
Ok(())
}
async fn ensure_projected_mcp_server_names_available(
&self,
projection: &crate::capability::CapabilityProjection,
) -> std::result::Result<(), crate::capability::CapabilityRuntimeError> {
let server_names = projection
.iter()
.filter_map(|(_, value)| match value {
crate::capability::CapabilityValue::Mcp(binding) => {
Some(binding.server_name().to_owned())
}
_ => None,
})
.collect::<Vec<_>>();
for server_name in server_names {
for manager in &self.mcp_managers {
if manager.contains_server(&server_name).await {
return Err(
crate::capability::CapabilityRuntimeError::RuntimeNameConflict {
kind: crate::capability::CapabilityKind::Mcp,
public_name: server_name,
},
);
}
}
}
Ok(())
}
}