1use super::*;
2
3impl AgentSession {
4 pub fn has_queue(&self) -> bool {
9 QueueControl::from_session(self).has_queue()
10 }
11
12 pub async fn set_lane_handler(
16 &self,
17 lane: SessionLane,
18 config: LaneHandlerConfig,
19 ) -> crate::error::Result<()> {
20 let _mutation = self.close_handle.extension_mutation.lock().await;
21 if self.is_closed() {
22 return Err(crate::error::CodeError::SessionClosed {
23 session_id: self.session_id.clone(),
24 });
25 }
26 QueueControl::from_session(self)
27 .set_lane_handler(lane, config)
28 .await;
29 if self.is_closed() {
30 return Err(crate::error::CodeError::SessionClosed {
31 session_id: self.session_id.clone(),
32 });
33 }
34 Ok(())
35 }
36
37 pub async fn complete_external_task(&self, task_id: &str, result: ExternalTaskResult) -> bool {
41 QueueControl::from_session(self)
42 .complete_external_task(task_id, result)
43 .await
44 }
45
46 pub async fn pending_external_tasks(&self) -> Vec<ExternalTask> {
48 QueueControl::from_session(self)
49 .pending_external_tasks()
50 .await
51 }
52
53 pub async fn queue_stats(&self) -> SessionQueueStats {
55 QueueControl::from_session(self).stats().await
56 }
57
58 pub async fn queue_metrics(&self) -> Option<MetricsSnapshot> {
60 QueueControl::from_session(self).metrics().await
61 }
62
63 pub async fn dead_letters(&self) -> Vec<DeadLetter> {
65 QueueControl::from_session(self).dead_letters().await
66 }
67
68 pub fn register_agent_dir(&self, dir: &std::path::Path) -> crate::error::Result<usize> {
81 let agents = crate::subagent::load_agents_from_dir(dir);
82 self.close_handle.mutate_immediate(|| {
83 for agent in &agents {
84 self.ensure_compatibility_name_available(
85 crate::capability::CapabilityKind::Agent,
86 &agent.name,
87 )?;
88 }
89 let count = agents.len();
90 for agent in agents {
91 tracing::info!(
92 session_id = %self.session_id,
93 agent = agent.name,
94 dir = %dir.display(),
95 "Dynamically registered agent"
96 );
97 self.agent_registry.register(agent);
98 }
99 Ok(count)
100 })?
101 }
102
103 pub fn register_worker_agent(
111 &self,
112 spec: crate::subagent::WorkerAgentSpec,
113 ) -> crate::error::Result<crate::subagent::AgentDefinition> {
114 self.close_handle.mutate_immediate(|| {
115 self.ensure_compatibility_name_available(
116 crate::capability::CapabilityKind::Agent,
117 &spec.name,
118 )?;
119 Ok(SessionExtensionRuntime::from_session(self).register_worker_agent(spec))
120 })?
121 }
122
123 pub fn register_worker_agents<I>(
125 &self,
126 specs: I,
127 ) -> crate::error::Result<Vec<crate::subagent::AgentDefinition>>
128 where
129 I: IntoIterator<Item = crate::subagent::WorkerAgentSpec>,
130 {
131 let specs = specs.into_iter().collect::<Vec<_>>();
132 self.close_handle.mutate_immediate(|| {
133 for spec in &specs {
134 self.ensure_compatibility_name_available(
135 crate::capability::CapabilityKind::Agent,
136 &spec.name,
137 )?;
138 }
139 Ok(SessionExtensionRuntime::from_session(self).register_worker_agents(specs))
140 })?
141 }
142
143 pub fn add_skill(&self, skill: Arc<crate::skills::Skill>) -> crate::error::Result<()> {
150 self.close_handle.mutate_immediate(|| {
151 self.ensure_compatibility_name_available(
152 crate::capability::CapabilityKind::Skill,
153 &skill.name,
154 )?;
155 SessionExtensionRuntime::from_session(self).add_skill(skill)
156 })?
157 }
158
159 pub fn remove_skill(&self, name: &str) -> crate::error::Result<()> {
164 self.close_handle
165 .mutate_immediate(|| SessionExtensionRuntime::from_session(self).remove_skill(name))
166 }
167
168 pub fn skill_names(&self) -> Vec<String> {
170 self.close_handle.skill_registry.list()
171 }
172
173 pub async fn add_mcp_server(
180 &self,
181 config: crate::mcp::McpServerConfig,
182 ) -> crate::error::Result<usize> {
183 SessionExtensionRuntime::from_session(self)
184 .add_mcp_server(config)
185 .await
186 }
187
188 #[cfg(feature = "serve")]
193 pub(crate) fn tool_executor(&self) -> &Arc<crate::tools::ToolExecutor> {
194 &self.tool_executor
195 }
196
197 pub fn register_dynamic_tool(
203 &self,
204 tool: Arc<dyn crate::tools::Tool>,
205 ) -> crate::error::Result<()> {
206 self.close_handle.mutate_immediate(|| {
207 self.ensure_compatibility_name_available(
208 crate::capability::CapabilityKind::Tool,
209 tool.name(),
210 )?;
211 self.tool_executor.register_dynamic_tool(tool);
212 Ok(())
213 })?
214 }
215
216 pub fn register_dynamic_workflow_runtime(&self) -> crate::error::Result<()> {
223 self.close_handle.mutate_immediate(|| {
224 self.ensure_compatibility_name_available(
225 crate::capability::CapabilityKind::Tool,
226 "dynamic_workflow",
227 )?;
228 crate::tools::register_dynamic_workflow(self.tool_executor.registry());
229 Ok(())
230 })?
231 }
232
233 pub fn unregister_dynamic_tool(&self, name: &str) -> crate::error::Result<()> {
236 self.close_handle
237 .mutate_immediate(|| self.tool_executor.unregister_dynamic_tool(name))
238 }
239
240 pub async fn remove_mcp_server(&self, server_name: &str) -> crate::error::Result<()> {
245 SessionExtensionRuntime::from_session(self)
246 .remove_mcp_server(server_name)
247 .await
248 }
249
250 pub async fn mcp_status(
252 &self,
253 ) -> std::collections::HashMap<String, crate::mcp::McpServerStatus> {
254 SessionExtensionRuntime::from_session(self)
255 .mcp_status()
256 .await
257 }
258
259 pub fn capability_catalog_stamp(&self) -> crate::capability::CapabilityCatalogStamp {
261 self.capability_catalog.current_stamp()
262 }
263
264 pub fn ensure_recovery_capability_binding(
267 &self,
268 expected: &crate::capability::RunCapabilityBindingV1,
269 ) -> std::result::Result<(), crate::capability::RunCapabilityBindingError> {
270 super::agent_loop_runtime::validate_run_capability_binding(self, expected)
271 }
272
273 pub async fn apply_capability_batch(
280 &self,
281 batch: crate::capability::SessionCapabilityBatch,
282 cancellation: tokio_util::sync::CancellationToken,
283 ) -> std::result::Result<
284 crate::capability::CapabilityCommitReceipt,
285 crate::capability::CapabilityRuntimeError,
286 > {
287 let _mutation = self.close_handle.extension_mutation.lock().await;
288 if self.is_closed() {
289 return Err(crate::capability::CapabilityRuntimeError::SessionClosed);
290 }
291
292 let preparation_cancellation = tokio_util::sync::CancellationToken::new();
293 let prepared = tokio::select! {
294 biased;
295 _ = self.session_cancel.cancelled() => {
296 preparation_cancellation.cancel();
297 return Err(if self.is_closed() {
298 crate::capability::CapabilityRuntimeError::SessionClosed
299 } else {
300 crate::capability::CapabilityRuntimeError::Cancelled
301 });
302 }
303 _ = cancellation.cancelled() => {
304 preparation_cancellation.cancel();
305 return Err(crate::capability::CapabilityRuntimeError::Cancelled);
306 }
307 result = batch.prepare(
308 &self.capability_catalog,
309 preparation_cancellation.clone(),
310 ) => result?,
311 };
312 self.ensure_projected_mcp_server_names_available(prepared.projection()?)
313 .await?;
314
315 let _publication = self
318 .close_handle
319 .immediate_extension_mutation
320 .lock()
321 .unwrap_or_else(std::sync::PoisonError::into_inner);
322 if self.is_closed() {
323 return Err(crate::capability::CapabilityRuntimeError::SessionClosed);
324 }
325 if cancellation.is_cancelled() || self.session_cancel.is_cancelled() {
326 return Err(crate::capability::CapabilityRuntimeError::Cancelled);
327 }
328 let command_registry = self
332 .command_registry
333 .lock()
334 .unwrap_or_else(std::sync::PoisonError::into_inner);
335 let current_projection = self.capability_catalog.pin();
336 super::agent_loop_runtime::validate_capability_projection_runtime(
337 self,
338 prepared.projection()?,
339 &command_registry,
340 )?;
341 super::agent_loop_runtime::validate_capability_projection_transition(
342 self,
343 current_projection.projection(),
344 prepared.projection()?,
345 )?;
346 prepared.commit()
347 }
348
349 pub async fn apply_sdk_capability_batch(
355 &self,
356 batch: crate::capability::SdkCapabilityBatchV1,
357 ) -> std::result::Result<
358 crate::capability::SdkCapabilityCommitReceiptV1,
359 crate::capability::SdkCapabilityBatchError,
360 > {
361 let session_batch = batch.into_session_batch()?;
362 let receipt = self
363 .apply_capability_batch(session_batch, self.session_cancel.child_token())
364 .await?;
365 Ok(crate::capability::SdkCapabilityCommitReceiptV1::from_receipt(&receipt))
366 }
367
368 pub async fn bootstrap_recovery_capability_batch(
378 &self,
379 expected: &crate::capability::RunCapabilityBindingV1,
380 batch: crate::capability::SessionCapabilityBatch,
381 cancellation: tokio_util::sync::CancellationToken,
382 ) -> std::result::Result<
383 crate::capability::CapabilityCommitReceipt,
384 crate::capability::CapabilityRuntimeError,
385 > {
386 let _mutation = self.close_handle.extension_mutation.lock().await;
387 if self.is_closed() {
388 return Err(crate::capability::CapabilityRuntimeError::SessionClosed);
389 }
390
391 let target_ceiling = self.capability_run_ceiling(batch.target())?;
392 expected
393 .ensure_matches(batch.target(), &target_ceiling)
394 .map_err(
395 |error| crate::capability::CapabilityRuntimeError::RecoveryBinding {
396 message: error.to_string(),
397 },
398 )?;
399
400 let preparation_cancellation = tokio_util::sync::CancellationToken::new();
401 let prepared = tokio::select! {
402 biased;
403 _ = self.session_cancel.cancelled() => {
404 preparation_cancellation.cancel();
405 return Err(if self.is_closed() {
406 crate::capability::CapabilityRuntimeError::SessionClosed
407 } else {
408 crate::capability::CapabilityRuntimeError::Cancelled
409 });
410 }
411 _ = cancellation.cancelled() => {
412 preparation_cancellation.cancel();
413 return Err(crate::capability::CapabilityRuntimeError::Cancelled);
414 }
415 result = batch.prepare_recovery_bootstrap(
416 &self.capability_catalog,
417 preparation_cancellation.clone(),
418 ) => result?,
419 };
420 self.ensure_projected_mcp_server_names_available(prepared.projection()?)
421 .await?;
422
423 let _publication = self
424 .close_handle
425 .immediate_extension_mutation
426 .lock()
427 .unwrap_or_else(std::sync::PoisonError::into_inner);
428 if self.is_closed() {
429 return Err(crate::capability::CapabilityRuntimeError::SessionClosed);
430 }
431 if cancellation.is_cancelled() || self.session_cancel.is_cancelled() {
432 return Err(crate::capability::CapabilityRuntimeError::Cancelled);
433 }
434 let command_registry = self
435 .command_registry
436 .lock()
437 .unwrap_or_else(std::sync::PoisonError::into_inner);
438 let current_projection = self.capability_catalog.pin();
439 super::agent_loop_runtime::validate_capability_projection_runtime(
440 self,
441 prepared.projection()?,
442 &command_registry,
443 )?;
444 super::agent_loop_runtime::validate_capability_projection_transition(
445 self,
446 current_projection.projection(),
447 prepared.projection()?,
448 )?;
449 prepared.commit()
450 }
451
452 pub async fn drain_capability_cleanup(&self) -> crate::capability::CapabilityCleanupReport {
454 self.capability_catalog.drain_cleanup().await
455 }
456
457 #[cfg(test)]
458 pub(crate) async fn admit_capability_run(
459 &self,
460 ) -> std::result::Result<
461 crate::capability::SessionCapabilityRun,
462 crate::capability::CapabilityRuntimeError,
463 > {
464 let projection = {
465 let _admission = self
466 .close_handle
467 .immediate_extension_mutation
468 .lock()
469 .unwrap_or_else(std::sync::PoisonError::into_inner);
470 if self.is_closed() {
471 return Err(crate::capability::CapabilityRuntimeError::SessionClosed);
472 }
473 self.capability_catalog.pin()
474 };
475 let ceiling = self.capability_run_ceiling(projection.projection().set())?;
476 crate::capability::SessionCapabilityRun::admit(
477 projection,
478 "active",
479 "active",
480 ceiling,
481 self.session_cancel.child_token(),
482 )
483 .await
484 }
485
486 pub(super) fn capability_run_ceiling(
487 &self,
488 set: &crate::capability::CapabilitySet,
489 ) -> std::result::Result<
490 crate::capability::CapabilityCeiling,
491 crate::capability::CapabilityRuntimeError,
492 > {
493 let mut governance = crate::capability::GovernanceCapabilityCeiling::none_required();
494 if self.config.permission_checker.is_some() || self.config.permission_policy.is_some() {
495 governance = governance.require_permission_guard();
496 }
497 if self.config.confirmation_manager.is_some() || self.config.confirmation_policy.is_some() {
498 governance = governance.require_confirmation_guard();
499 }
500 if self.config.security_provider.is_some() {
501 governance = governance.require_security_guard();
502 }
503 if self.config.budget_guard.is_some() || self.budget_guard().is_some() {
504 governance = governance.require_budget_guard();
505 }
506 if self.config.enforce_active_skill_tool_restrictions {
507 governance = governance.require_active_skill_restrictions();
508 }
509 let execution = crate::capability::CapabilityExecutionCeiling::new(
510 self.config.max_tool_rounds,
511 self.config.max_parallel_tasks,
512 self.config.tool_timeout_ms,
513 self.config.llm_api_timeout_ms,
514 self.config.max_execution_time_ms,
515 )?;
516 crate::capability::CapabilityCeiling::all(
517 set,
518 crate::capability::WorkspaceCapabilityCeiling::all(),
519 governance,
520 execution,
521 )
522 .map_err(Into::into)
523 }
524
525 pub(super) fn ensure_compatibility_name_available(
526 &self,
527 kind: crate::capability::CapabilityKind,
528 public_name: &str,
529 ) -> crate::error::Result<()> {
530 let projection = self.capability_catalog.pin();
531 if projection
532 .projection()
533 .iter()
534 .any(|(_, value)| match value {
535 crate::capability::CapabilityValue::Mcp(binding)
536 if kind == crate::capability::CapabilityKind::Tool =>
537 {
538 binding.contains_public_tool_name(public_name)
539 }
540 crate::capability::CapabilityValue::Agent(agent) => {
541 kind == crate::capability::CapabilityKind::Agent
542 && crate::subagent::agent_names_conflict(&agent.name, public_name)
543 }
544 _ => value.kind() == kind && value.public_name() == Some(public_name),
545 })
546 {
547 return Err(
548 crate::capability::CapabilityRuntimeError::RuntimeNameConflict {
549 kind,
550 public_name: public_name.to_owned(),
551 }
552 .into(),
553 );
554 }
555 Ok(())
556 }
557
558 async fn ensure_projected_mcp_server_names_available(
559 &self,
560 projection: &crate::capability::CapabilityProjection,
561 ) -> std::result::Result<(), crate::capability::CapabilityRuntimeError> {
562 let server_names = projection
563 .iter()
564 .filter_map(|(_, value)| match value {
565 crate::capability::CapabilityValue::Mcp(binding) => {
566 Some(binding.server_name().to_owned())
567 }
568 _ => None,
569 })
570 .collect::<Vec<_>>();
571 for server_name in server_names {
572 for manager in &self.mcp_managers {
573 if manager.contains_server(&server_name).await {
574 return Err(
575 crate::capability::CapabilityRuntimeError::RuntimeNameConflict {
576 kind: crate::capability::CapabilityKind::Mcp,
577 public_name: server_name,
578 },
579 );
580 }
581 }
582 }
583 Ok(())
584 }
585}