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 pub fn register_dynamic_tool(
194 &self,
195 tool: Arc<dyn crate::tools::Tool>,
196 ) -> crate::error::Result<()> {
197 self.close_handle.mutate_immediate(|| {
198 self.ensure_compatibility_name_available(
199 crate::capability::CapabilityKind::Tool,
200 tool.name(),
201 )?;
202 self.tool_executor.register_dynamic_tool(tool);
203 Ok(())
204 })?
205 }
206
207 #[cfg(feature = "dynamic-workflow")]
214 pub fn register_dynamic_workflow_runtime(&self) -> crate::error::Result<()> {
215 self.close_handle.mutate_immediate(|| {
216 self.ensure_compatibility_name_available(
217 crate::capability::CapabilityKind::Tool,
218 "dynamic_workflow",
219 )?;
220 crate::tools::register_dynamic_workflow(self.tool_executor.registry());
221 Ok(())
222 })?
223 }
224
225 #[cfg(not(feature = "dynamic-workflow"))]
228 pub fn register_dynamic_workflow_runtime(&self) -> crate::error::Result<()> {
229 self.close_handle.mutate_immediate(|| {
230 Err(crate::error::CodeError::Config(
231 "dynamic_workflow requires the advanced-harness / dynamic-workflow feature".into(),
232 ))
233 })?
234 }
235
236 pub fn unregister_dynamic_tool(&self, name: &str) -> crate::error::Result<()> {
239 self.close_handle
240 .mutate_immediate(|| self.tool_executor.unregister_dynamic_tool(name))
241 }
242
243 pub async fn remove_mcp_server(&self, server_name: &str) -> crate::error::Result<()> {
248 SessionExtensionRuntime::from_session(self)
249 .remove_mcp_server(server_name)
250 .await
251 }
252
253 pub async fn republish_inherited_mcp_tools(&self) -> crate::error::Result<()> {
260 SessionExtensionRuntime::from_session(self)
261 .republish_inherited_mcp_tools()
262 .await
263 }
264
265 pub fn inherits_mcp_managers(&self) -> bool {
268 !self.inherited_mcp_managers.is_empty()
269 }
270
271 pub async fn mcp_status(
273 &self,
274 ) -> std::collections::HashMap<String, crate::mcp::McpServerStatus> {
275 SessionExtensionRuntime::from_session(self)
276 .mcp_status()
277 .await
278 }
279
280 pub fn capability_catalog_stamp(&self) -> crate::capability::CapabilityCatalogStamp {
282 self.capability_catalog.current_stamp()
283 }
284
285 pub fn ensure_recovery_capability_binding(
288 &self,
289 expected: &crate::capability::RunCapabilityBindingV1,
290 ) -> std::result::Result<(), crate::capability::RunCapabilityBindingError> {
291 super::agent_loop_runtime::validate_run_capability_binding(self, expected)
292 }
293
294 pub async fn apply_capability_batch(
301 &self,
302 batch: crate::capability::SessionCapabilityBatch,
303 cancellation: tokio_util::sync::CancellationToken,
304 ) -> std::result::Result<
305 crate::capability::CapabilityCommitReceipt,
306 crate::capability::CapabilityRuntimeError,
307 > {
308 let _mutation = self.close_handle.extension_mutation.lock().await;
309 if self.is_closed() {
310 return Err(crate::capability::CapabilityRuntimeError::SessionClosed);
311 }
312
313 let preparation_cancellation = tokio_util::sync::CancellationToken::new();
314 let prepared = tokio::select! {
315 biased;
316 _ = self.session_cancel.cancelled() => {
317 preparation_cancellation.cancel();
318 return Err(if self.is_closed() {
319 crate::capability::CapabilityRuntimeError::SessionClosed
320 } else {
321 crate::capability::CapabilityRuntimeError::Cancelled
322 });
323 }
324 _ = cancellation.cancelled() => {
325 preparation_cancellation.cancel();
326 return Err(crate::capability::CapabilityRuntimeError::Cancelled);
327 }
328 result = batch.prepare(
329 &self.capability_catalog,
330 preparation_cancellation.clone(),
331 ) => result?,
332 };
333 self.ensure_projected_mcp_server_names_available(prepared.projection()?)
334 .await?;
335
336 let _publication = self
339 .close_handle
340 .immediate_extension_mutation
341 .lock()
342 .unwrap_or_else(std::sync::PoisonError::into_inner);
343 if self.is_closed() {
344 return Err(crate::capability::CapabilityRuntimeError::SessionClosed);
345 }
346 if cancellation.is_cancelled() || self.session_cancel.is_cancelled() {
347 return Err(crate::capability::CapabilityRuntimeError::Cancelled);
348 }
349 let command_registry = self
353 .command_registry
354 .lock()
355 .unwrap_or_else(std::sync::PoisonError::into_inner);
356 let current_projection = self.capability_catalog.pin();
357 super::agent_loop_runtime::validate_capability_projection_runtime(
358 self,
359 prepared.projection()?,
360 &command_registry,
361 )?;
362 super::agent_loop_runtime::validate_capability_projection_transition(
363 self,
364 current_projection.projection(),
365 prepared.projection()?,
366 )?;
367 prepared.commit()
368 }
369
370 pub async fn apply_sdk_capability_batch(
376 &self,
377 batch: crate::capability::SdkCapabilityBatchV1,
378 ) -> std::result::Result<
379 crate::capability::SdkCapabilityCommitReceiptV1,
380 crate::capability::SdkCapabilityBatchError,
381 > {
382 let session_batch = batch.into_session_batch()?;
383 let receipt = self
384 .apply_capability_batch(session_batch, self.session_cancel.child_token())
385 .await?;
386 Ok(crate::capability::SdkCapabilityCommitReceiptV1::from_receipt(&receipt))
387 }
388
389 pub async fn bootstrap_recovery_capability_batch(
399 &self,
400 expected: &crate::capability::RunCapabilityBindingV1,
401 batch: crate::capability::SessionCapabilityBatch,
402 cancellation: tokio_util::sync::CancellationToken,
403 ) -> std::result::Result<
404 crate::capability::CapabilityCommitReceipt,
405 crate::capability::CapabilityRuntimeError,
406 > {
407 let _mutation = self.close_handle.extension_mutation.lock().await;
408 if self.is_closed() {
409 return Err(crate::capability::CapabilityRuntimeError::SessionClosed);
410 }
411
412 let target_ceiling = self.capability_run_ceiling(batch.target())?;
413 expected
414 .ensure_matches(batch.target(), &target_ceiling)
415 .map_err(
416 |error| crate::capability::CapabilityRuntimeError::RecoveryBinding {
417 message: error.to_string(),
418 },
419 )?;
420
421 let preparation_cancellation = tokio_util::sync::CancellationToken::new();
422 let prepared = tokio::select! {
423 biased;
424 _ = self.session_cancel.cancelled() => {
425 preparation_cancellation.cancel();
426 return Err(if self.is_closed() {
427 crate::capability::CapabilityRuntimeError::SessionClosed
428 } else {
429 crate::capability::CapabilityRuntimeError::Cancelled
430 });
431 }
432 _ = cancellation.cancelled() => {
433 preparation_cancellation.cancel();
434 return Err(crate::capability::CapabilityRuntimeError::Cancelled);
435 }
436 result = batch.prepare_recovery_bootstrap(
437 &self.capability_catalog,
438 preparation_cancellation.clone(),
439 ) => result?,
440 };
441 self.ensure_projected_mcp_server_names_available(prepared.projection()?)
442 .await?;
443
444 let _publication = self
445 .close_handle
446 .immediate_extension_mutation
447 .lock()
448 .unwrap_or_else(std::sync::PoisonError::into_inner);
449 if self.is_closed() {
450 return Err(crate::capability::CapabilityRuntimeError::SessionClosed);
451 }
452 if cancellation.is_cancelled() || self.session_cancel.is_cancelled() {
453 return Err(crate::capability::CapabilityRuntimeError::Cancelled);
454 }
455 let command_registry = self
456 .command_registry
457 .lock()
458 .unwrap_or_else(std::sync::PoisonError::into_inner);
459 let current_projection = self.capability_catalog.pin();
460 super::agent_loop_runtime::validate_capability_projection_runtime(
461 self,
462 prepared.projection()?,
463 &command_registry,
464 )?;
465 super::agent_loop_runtime::validate_capability_projection_transition(
466 self,
467 current_projection.projection(),
468 prepared.projection()?,
469 )?;
470 prepared.commit()
471 }
472
473 pub async fn drain_capability_cleanup(&self) -> crate::capability::CapabilityCleanupReport {
475 self.capability_catalog.drain_cleanup().await
476 }
477
478 #[cfg(test)]
479 pub(crate) async fn admit_capability_run(
480 &self,
481 ) -> std::result::Result<
482 crate::capability::SessionCapabilityRun,
483 crate::capability::CapabilityRuntimeError,
484 > {
485 let projection = {
486 let _admission = self
487 .close_handle
488 .immediate_extension_mutation
489 .lock()
490 .unwrap_or_else(std::sync::PoisonError::into_inner);
491 if self.is_closed() {
492 return Err(crate::capability::CapabilityRuntimeError::SessionClosed);
493 }
494 self.capability_catalog.pin()
495 };
496 let ceiling = self.capability_run_ceiling(projection.projection().set())?;
497 crate::capability::SessionCapabilityRun::admit(
498 projection,
499 "active",
500 "active",
501 ceiling,
502 self.session_cancel.child_token(),
503 )
504 .await
505 }
506
507 pub(super) fn capability_run_ceiling(
508 &self,
509 set: &crate::capability::CapabilitySet,
510 ) -> std::result::Result<
511 crate::capability::CapabilityCeiling,
512 crate::capability::CapabilityRuntimeError,
513 > {
514 let mut governance = crate::capability::GovernanceCapabilityCeiling::none_required();
515 if self.config.permission_checker.is_some() || self.config.permission_policy.is_some() {
516 governance = governance.require_permission_guard();
517 }
518 if self.config.confirmation_manager.is_some() || self.config.confirmation_policy.is_some() {
519 governance = governance.require_confirmation_guard();
520 }
521 if self.config.security_provider.is_some() {
522 governance = governance.require_security_guard();
523 }
524 if self.config.budget_guard.is_some() || self.budget_guard().is_some() {
525 governance = governance.require_budget_guard();
526 }
527 if self.config.enforce_active_skill_tool_restrictions {
528 governance = governance.require_active_skill_restrictions();
529 }
530 let execution = crate::capability::CapabilityExecutionCeiling::new(
531 self.config.max_tool_rounds,
532 self.config.max_parallel_tasks,
533 self.config.tool_timeout_ms,
534 self.config.llm_api_timeout_ms,
535 self.config.max_execution_time_ms,
536 )?;
537 crate::capability::CapabilityCeiling::all(
538 set,
539 crate::capability::WorkspaceCapabilityCeiling::all(),
540 governance,
541 execution,
542 )
543 .map_err(Into::into)
544 }
545
546 pub(super) fn ensure_compatibility_name_available(
547 &self,
548 kind: crate::capability::CapabilityKind,
549 public_name: &str,
550 ) -> crate::error::Result<()> {
551 let projection = self.capability_catalog.pin();
552 if projection
553 .projection()
554 .iter()
555 .any(|(_, value)| match value {
556 crate::capability::CapabilityValue::Mcp(binding)
557 if kind == crate::capability::CapabilityKind::Tool =>
558 {
559 binding.contains_public_tool_name(public_name)
560 }
561 crate::capability::CapabilityValue::Agent(agent) => {
562 kind == crate::capability::CapabilityKind::Agent
563 && crate::subagent::agent_names_conflict(&agent.name, public_name)
564 }
565 _ => value.kind() == kind && value.public_name() == Some(public_name),
566 })
567 {
568 return Err(
569 crate::capability::CapabilityRuntimeError::RuntimeNameConflict {
570 kind,
571 public_name: public_name.to_owned(),
572 }
573 .into(),
574 );
575 }
576 Ok(())
577 }
578
579 async fn ensure_projected_mcp_server_names_available(
580 &self,
581 projection: &crate::capability::CapabilityProjection,
582 ) -> std::result::Result<(), crate::capability::CapabilityRuntimeError> {
583 let server_names = projection
584 .iter()
585 .filter_map(|(_, value)| match value {
586 crate::capability::CapabilityValue::Mcp(binding) => {
587 Some(binding.server_name().to_owned())
588 }
589 _ => None,
590 })
591 .collect::<Vec<_>>();
592 for server_name in server_names {
593 for manager in &self.mcp_managers {
594 if manager.contains_server(&server_name).await {
595 return Err(
596 crate::capability::CapabilityRuntimeError::RuntimeNameConflict {
597 kind: crate::capability::CapabilityKind::Mcp,
598 public_name: server_name,
599 },
600 );
601 }
602 }
603 }
604 Ok(())
605 }
606}