use crate::agent::{
AgentToolDispatcher, BindOutcome, DispatcherCapabilities, ExternalToolUpdate,
OpsLifecycleBindError, ToolDispatchContext,
};
use crate::error::ToolError;
use crate::ops::ToolAccessPolicy;
use crate::tool_catalog::{ToolCatalogCapabilities, ToolCatalogEntry};
use crate::types::{ToolCallView, ToolDef, ToolNameSet};
use async_trait::async_trait;
use std::sync::Arc;
#[cfg_attr(target_arch = "wasm32", async_trait(?Send))]
#[cfg_attr(not(target_arch = "wasm32"), async_trait)]
pub trait ToolDispatchAdmission: Send + Sync {
async fn await_dispatch_admission(
&self,
call: ToolCallView<'_>,
context: Option<&ToolDispatchContext>,
effect_kind: crate::LiveBridgeEffectKind,
) -> Result<(), ToolError>;
async fn record_dispatch_outcome(
&self,
_call: ToolCallView<'_>,
_context: Option<&ToolDispatchContext>,
_effect_kind: crate::LiveBridgeEffectKind,
_outcome: crate::LiveBridgeEffectOutcome,
) -> Result<(), ToolError> {
Ok(())
}
}
#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
pub enum ToolExecutionPolicyError {
#[error(
"tool access policy 'inherit' is unresolved at the execution seam; \
the spawn chain must resolve it to the parent's effective policy \
before the dispatch gate is built"
)]
UnresolvedInherit,
#[error("tool access constraints must not be empty")]
EmptyConstraints,
}
impl ToolExecutionPolicyError {
pub fn error_code(&self) -> &'static str {
match self {
Self::UnresolvedInherit => "tool_execution_policy_unresolved_inherit",
Self::EmptyConstraints => "tool_execution_policy_empty_constraints",
}
}
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
pub enum ToolMutationClass {
ReadOnly,
Mutating,
#[default]
Unknown,
}
impl ToolMutationClass {
#[must_use]
pub const fn is_declared_read_only(self) -> bool {
matches!(self, Self::ReadOnly)
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ToolExecutionPolicy {
constraints: Vec<crate::ops::ToolAccessConstraint>,
}
impl ToolExecutionPolicy {
#[must_use]
pub fn unrestricted() -> Self {
Self {
constraints: Vec::new(),
}
}
pub fn resolve(policy: ToolAccessPolicy) -> Result<Self, ToolExecutionPolicyError> {
let constraints = match policy {
ToolAccessPolicy::Inherit => Err(ToolExecutionPolicyError::UnresolvedInherit),
ToolAccessPolicy::AllowList(names) => {
Ok(vec![crate::ops::ToolAccessConstraint::AllowNames(names)])
}
ToolAccessPolicy::DenyList(names) => {
Ok(vec![crate::ops::ToolAccessConstraint::DenyNames(names)])
}
ToolAccessPolicy::ReadOnly => Ok(vec![crate::ops::ToolAccessConstraint::ReadOnly]),
ToolAccessPolicy::Constraints(constraints) if constraints.is_empty() => {
Err(ToolExecutionPolicyError::EmptyConstraints)
}
ToolAccessPolicy::Constraints(constraints) => Ok(constraints),
}?;
Ok(Self {
constraints: normalize_constraints(constraints),
})
}
#[must_use]
pub fn is_unrestricted(&self) -> bool {
self.constraints.is_empty()
}
#[must_use]
pub fn is_read_only_intent(&self) -> bool {
self.constraints.len() == 1
&& matches!(
self.constraints[0],
crate::ops::ToolAccessConstraint::ReadOnly
)
}
#[must_use]
fn requires_mutation_declaration(&self) -> bool {
self.constraints
.iter()
.any(|constraint| matches!(constraint, crate::ops::ToolAccessConstraint::ReadOnly))
}
#[must_use]
pub fn permits_call(&self, name: &str, declared: ToolMutationClass) -> bool {
self.constraints.iter().all(|constraint| match constraint {
crate::ops::ToolAccessConstraint::AllowNames(names) => names.contains(name),
crate::ops::ToolAccessConstraint::DenyNames(names) => !names.contains(name),
crate::ops::ToolAccessConstraint::ReadOnly => declared.is_declared_read_only(),
})
}
#[must_use]
pub fn permits(&self, name: &str) -> bool {
self.permits_call(name, ToolMutationClass::Unknown)
}
}
fn normalize_constraints(
constraints: Vec<crate::ops::ToolAccessConstraint>,
) -> Vec<crate::ops::ToolAccessConstraint> {
use crate::ops::ToolAccessConstraint;
let mut allow: Option<ToolNameSet> = None;
let mut deny = ToolNameSet::default();
let mut read_only = false;
for constraint in constraints {
match constraint {
ToolAccessConstraint::AllowNames(names) => {
allow = Some(match allow {
None => names,
Some(mut existing) => {
existing.retain(|name| names.contains(name.as_str()));
existing
}
});
}
ToolAccessConstraint::DenyNames(names) => {
for name in names.into_inner() {
deny.insert(name);
}
}
ToolAccessConstraint::ReadOnly => read_only = true,
}
}
let mut normalized = Vec::new();
if let Some(names) = allow {
normalized.push(ToolAccessConstraint::AllowNames(names));
}
if !deny.is_empty() {
normalized.push(ToolAccessConstraint::DenyNames(deny));
}
if read_only {
normalized.push(ToolAccessConstraint::ReadOnly);
}
normalized
}
pub struct ExecutionPolicyGatedDispatcher<T: AgentToolDispatcher + ?Sized> {
inner: Arc<T>,
policy: ToolExecutionPolicy,
consequence_policy: Option<crate::BoundToolConsequencePolicy>,
dispatch_admission: Option<Arc<dyn ToolDispatchAdmission>>,
}
impl<T: AgentToolDispatcher + ?Sized> ExecutionPolicyGatedDispatcher<T> {
pub fn new(inner: Arc<T>, policy: ToolExecutionPolicy) -> Self {
Self {
inner,
policy,
consequence_policy: None,
dispatch_admission: None,
}
}
#[must_use]
pub fn with_consequence_policy(
mut self,
consequence_policy: crate::BoundToolConsequencePolicy,
) -> Self {
self.consequence_policy = Some(consequence_policy);
self
}
#[must_use]
pub fn with_dispatch_admission(mut self, admission: Arc<dyn ToolDispatchAdmission>) -> Self {
self.dispatch_admission = Some(admission);
self
}
fn permits_inner_call(&self, name: &str) -> bool {
if self.policy.requires_mutation_declaration() {
return self
.policy
.permits_call(name, self.inner.tool_mutation_class(name));
}
self.policy.permits(name)
}
fn denial_error(&self, name: &str) -> ToolError {
let inner_knows_tool = if self.inner.tool_catalog_capabilities().exact_catalog {
self.inner
.tool_catalog()
.iter()
.any(|entry| entry.tool.name == name)
} else {
self.inner.tools().iter().any(|tool| tool.name == name)
};
if inner_knows_tool {
ToolError::access_denied(name)
} else {
ToolError::not_found(name)
}
}
async fn evaluate_consequence_policy(
&self,
call: ToolCallView<'_>,
context: Option<&ToolDispatchContext>,
) -> Result<(), ToolError> {
let Some(policy) = self.consequence_policy.as_ref() else {
return Ok(());
};
policy
.evaluate(call, context.and_then(ToolDispatchContext::run_id).cloned())
.await
}
async fn await_dispatch_admission(
&self,
call: ToolCallView<'_>,
context: Option<&ToolDispatchContext>,
) -> Result<DispatchAdmissionCustody, ToolError> {
let effect_kind = self.inner.live_bridge_effect_kind(call.name);
let mut admissions = Vec::with_capacity(2);
if let Some(admission) = self.dispatch_admission.as_ref() {
admission
.await_dispatch_admission(call, context, effect_kind)
.await?;
admissions.push(Arc::clone(admission));
}
if let Some(admission) = context.and_then(ToolDispatchContext::live_bridge_admission) {
if let Err(error) = admission
.admission()
.await_dispatch_admission(call, context, effect_kind)
.await
{
for admitted in admissions {
admitted
.record_dispatch_outcome(
call,
context,
effect_kind,
crate::LiveBridgeEffectOutcome::Failed,
)
.await?;
}
return Err(error);
}
admissions.push(Arc::clone(admission.admission()));
}
Ok(DispatchAdmissionCustody {
admissions,
effect_kind,
})
}
}
struct DispatchAdmissionCustody {
admissions: Vec<Arc<dyn ToolDispatchAdmission>>,
effect_kind: crate::LiveBridgeEffectKind,
}
impl DispatchAdmissionCustody {
async fn settle(
self,
call: ToolCallView<'_>,
context: Option<&ToolDispatchContext>,
outcome: crate::LiveBridgeEffectOutcome,
) -> Result<(), ToolError> {
for admission in self.admissions {
admission
.record_dispatch_outcome(call, context, self.effect_kind, outcome)
.await?;
}
Ok(())
}
}
#[cfg_attr(target_arch = "wasm32", async_trait(?Send))]
#[cfg_attr(not(target_arch = "wasm32"), async_trait)]
impl<T: AgentToolDispatcher + ?Sized + 'static> AgentToolDispatcher
for ExecutionPolicyGatedDispatcher<T>
{
fn tools(&self) -> Arc<[Arc<ToolDef>]> {
self.inner.tools()
}
fn tool_catalog_capabilities(&self) -> ToolCatalogCapabilities {
self.inner.tool_catalog_capabilities()
}
fn tool_catalog(&self) -> Arc<[ToolCatalogEntry]> {
self.inner.tool_catalog()
}
fn tool_mutation_class(&self, tool_name: &str) -> ToolMutationClass {
self.inner.tool_mutation_class(tool_name)
}
fn live_bridge_effect_kind(&self, tool_name: &str) -> crate::LiveBridgeEffectKind {
self.inner.live_bridge_effect_kind(tool_name)
}
fn execution_binding_epoch(&self, tool_name: &str) -> u64 {
self.inner.execution_binding_epoch(tool_name)
}
fn pending_catalog_sources(&self) -> Arc<[String]> {
self.inner.pending_catalog_sources()
}
fn execution_binding_fingerprint(
&self,
tool_name: &str,
) -> Result<crate::EphemeralToolBindingFingerprint, crate::ToolExecutionResolutionError> {
if !self.permits_inner_call(tool_name) {
return Err(match self.denial_error(tool_name) {
ToolError::NotFound { .. } => crate::ToolExecutionResolutionError::NotFound {
tool_name: tool_name.to_string(),
},
_ => crate::ToolExecutionResolutionError::AccessDenied {
tool_name: tool_name.to_string(),
},
});
}
let catalog = self.tool_catalog();
let entry = catalog
.iter()
.find(|entry| entry.tool.name == tool_name)
.ok_or_else(|| crate::ToolExecutionResolutionError::NotFound {
tool_name: tool_name.to_string(),
})?;
Ok(crate::ephemeral_tool_catalog_binding_fingerprint(entry)
.with_live_authority(0, 0)
.with_dependency(&self.inner.execution_binding_fingerprint(tool_name)?))
}
fn resolve_execution_plan(
&self,
call: ToolCallView<'_>,
dispatch_context: &ToolDispatchContext,
resolution_context: &crate::ToolExecutionResolutionContext,
) -> Result<crate::ResolvedToolExecutionPlan, crate::ToolExecutionResolutionError> {
if !self.permits_inner_call(call.name) {
return Err(match self.denial_error(call.name) {
ToolError::NotFound { .. } => crate::ToolExecutionResolutionError::NotFound {
tool_name: call.name.to_string(),
},
_ => crate::ToolExecutionResolutionError::AccessDenied {
tool_name: call.name.to_string(),
},
});
}
self.inner
.resolve_execution_plan(call, dispatch_context, resolution_context)
}
fn validate_resolved_execution_plan(
&self,
call: ToolCallView<'_>,
resolution_context: &crate::ToolExecutionResolutionContext,
plan: &crate::ResolvedToolExecutionPlan,
) -> Result<(), crate::ToolExecutionResolutionError> {
if !self.permits_inner_call(call.name) {
return Err(match self.denial_error(call.name) {
ToolError::NotFound { .. } => crate::ToolExecutionResolutionError::NotFound {
tool_name: call.name.to_string(),
},
_ => crate::ToolExecutionResolutionError::AccessDenied {
tool_name: call.name.to_string(),
},
});
}
self.inner
.validate_resolved_execution_plan(call, resolution_context, plan)
}
async fn dispatch(
&self,
call: ToolCallView<'_>,
) -> Result<crate::ops::ToolDispatchOutcome, ToolError> {
let custody = self.await_dispatch_admission(call, None).await?;
let pre_dispatch = async {
if !self.permits_inner_call(call.name) {
return Err(self.denial_error(call.name));
}
self.evaluate_consequence_policy(call, None).await
}
.await;
if let Err(error) = pre_dispatch {
custody
.settle(call, None, crate::LiveBridgeEffectOutcome::Failed)
.await?;
return Err(error);
}
let result = self.inner.dispatch(call).await;
let outcome = if result.is_ok() {
crate::LiveBridgeEffectOutcome::Committed
} else {
crate::LiveBridgeEffectOutcome::Unknown
};
custody.settle(call, None, outcome).await?;
result
}
async fn dispatch_with_context(
&self,
call: ToolCallView<'_>,
context: &ToolDispatchContext,
) -> Result<crate::ops::ToolDispatchOutcome, ToolError> {
let custody = self.await_dispatch_admission(call, Some(context)).await?;
let pre_dispatch = async {
if !self.permits_inner_call(call.name) {
return Err(self.denial_error(call.name));
}
self.evaluate_consequence_policy(call, Some(context)).await
}
.await;
if let Err(error) = pre_dispatch {
custody
.settle(call, Some(context), crate::LiveBridgeEffectOutcome::Failed)
.await?;
return Err(error);
}
let result = self.inner.dispatch_with_context(call, context).await;
let outcome = if result.is_ok() {
crate::LiveBridgeEffectOutcome::Committed
} else {
crate::LiveBridgeEffectOutcome::Unknown
};
custody.settle(call, Some(context), outcome).await?;
result
}
async fn dispatch_resolved_with_context(
&self,
call: ToolCallView<'_>,
context: &ToolDispatchContext,
plan: &crate::ResolvedToolExecutionPlan,
) -> Result<crate::ops::ToolDispatchOutcome, ToolError> {
let custody = self.await_dispatch_admission(call, Some(context)).await?;
let pre_dispatch = async {
if !self.permits_inner_call(call.name) {
return Err(self.denial_error(call.name));
}
self.evaluate_consequence_policy(call, Some(context)).await
}
.await;
if let Err(error) = pre_dispatch {
custody
.settle(call, Some(context), crate::LiveBridgeEffectOutcome::Failed)
.await?;
return Err(error);
}
let result = self
.inner
.dispatch_resolved_with_context(call, context, plan)
.await;
let outcome = if result.is_ok() {
crate::LiveBridgeEffectOutcome::Committed
} else {
crate::LiveBridgeEffectOutcome::Unknown
};
custody.settle(call, Some(context), outcome).await?;
result
}
async fn poll_external_updates(&self) -> ExternalToolUpdate {
self.inner.poll_external_updates().await
}
fn external_tool_surface_snapshot(&self) -> Option<crate::ExternalToolSurfaceSnapshot> {
self.inner.external_tool_surface_snapshot()
}
fn capabilities(&self) -> DispatcherCapabilities {
self.inner.capabilities()
}
fn bind_ops_lifecycle(
self: Arc<Self>,
registry: Arc<dyn crate::ops_lifecycle::OpsLifecycleRegistry>,
owner_bridge_session_id: crate::types::SessionId,
) -> Result<BindOutcome, OpsLifecycleBindError> {
let owned = Arc::try_unwrap(self).map_err(|_| OpsLifecycleBindError::SharedOwnership)?;
if Arc::strong_count(&owned.inner) == 1 {
let outcome = owned
.inner
.bind_ops_lifecycle(registry, owner_bridge_session_id)?;
let bound = outcome.was_bound();
let inner = outcome.into_dispatcher();
let gated = Arc::new(ExecutionPolicyGatedDispatcher {
inner,
policy: owned.policy,
consequence_policy: owned.consequence_policy,
dispatch_admission: owned.dispatch_admission,
});
Ok(if bound {
BindOutcome::Bound(gated)
} else {
BindOutcome::Skipped(gated)
})
} else {
Ok(BindOutcome::Skipped(Arc::new(
ExecutionPolicyGatedDispatcher {
inner: owned.inner,
policy: owned.policy,
consequence_policy: owned.consequence_policy,
dispatch_admission: owned.dispatch_admission,
},
)))
}
}
fn completion_enrichment(
&self,
) -> Option<Arc<dyn crate::completion_feed::CompletionEnrichmentProvider>> {
self.inner.completion_enrichment()
}
fn bind_mcp_server_lifecycle_handle(
&self,
handle: Arc<dyn crate::handles::McpServerLifecycleHandle>,
) {
self.inner.bind_mcp_server_lifecycle_handle(handle);
}
fn bind_external_tool_surface_handle(
&self,
handle: Arc<dyn crate::handles::ExternalToolSurfaceHandle>,
) {
self.inner.bind_external_tool_surface_handle(handle);
}
}
#[cfg(test)]
#[allow(clippy::unwrap_used, clippy::expect_used, clippy::panic)]
mod tests {
use super::*;
use crate::handles::{
DslTransitionError, ExternalToolSurfaceHandle, ExternalToolSurfaceInput,
ExternalToolSurfaceTransition, McpServerLifecycleHandle, SurfaceDiagnosticSnapshot,
SurfaceSnapshot,
};
use crate::ops_lifecycle::{
OperationCompletionWatch, OperationLifecycleSnapshot, OperationPeerHandle,
OperationProgressUpdate, OpsLifecycleError, OpsLifecycleRegistry,
};
use crate::tool_scope::ExternalToolSurfaceGlobalPhase;
use crate::types::ToolResult;
use crate::{
BoundToolConsequencePolicy, MobMemberBinding, PolicyDigest, PolicyEvaluationProvenance,
PolicyEvaluationSupervisorConfig, PolicyId, PolicyProviderGeneration, PolicyProviderId,
PolicyRevision, ToolConsequenceDenial, ToolConsequenceFailure,
ToolConsequenceNarrowingPolicy, ToolConsequencePolicyRegistry,
ToolConsequencePolicySnapshot, ToolConsequenceRequest, ToolConsequenceVerdict,
};
use std::collections::BTreeSet;
use std::sync::Mutex;
use std::sync::atomic::{AtomicBool, Ordering};
struct BlockingAdmission {
released: AtomicBool,
notify: tokio::sync::Notify,
effect_kinds: Mutex<Vec<crate::LiveBridgeEffectKind>>,
outcomes: Mutex<Vec<crate::LiveBridgeEffectOutcome>>,
}
impl BlockingAdmission {
fn new() -> Self {
Self {
released: AtomicBool::new(false),
notify: tokio::sync::Notify::new(),
effect_kinds: Mutex::new(Vec::new()),
outcomes: Mutex::new(Vec::new()),
}
}
fn release(&self) {
self.released.store(true, Ordering::Release);
self.notify.notify_waiters();
}
}
#[async_trait::async_trait]
impl ToolDispatchAdmission for BlockingAdmission {
async fn await_dispatch_admission(
&self,
_call: ToolCallView<'_>,
_context: Option<&ToolDispatchContext>,
effect_kind: crate::LiveBridgeEffectKind,
) -> Result<(), ToolError> {
self.effect_kinds.lock().unwrap().push(effect_kind);
loop {
let notified = self.notify.notified();
if self.released.load(Ordering::Acquire) {
return Ok(());
}
notified.await;
}
}
async fn record_dispatch_outcome(
&self,
_call: ToolCallView<'_>,
_context: Option<&ToolDispatchContext>,
_effect_kind: crate::LiveBridgeEffectKind,
outcome: crate::LiveBridgeEffectOutcome,
) -> Result<(), ToolError> {
self.outcomes.lock().unwrap().push(outcome);
Ok(())
}
}
fn tool_def(name: &str) -> Arc<ToolDef> {
Arc::new(ToolDef::new(
name,
format!("test tool {name}"),
serde_json::json!({ "type": "object" }),
))
}
fn empty_args() -> Box<serde_json::value::RawValue> {
serde_json::value::RawValue::from_string("{}".to_string()).expect("valid args json")
}
struct SpyDispatcher {
tools: Arc<[Arc<ToolDef>]>,
dispatched: Mutex<Vec<String>>,
ops_bound: Mutex<bool>,
mcp_handles_bound: Mutex<usize>,
surface_handles_bound: Mutex<usize>,
}
impl SpyDispatcher {
fn new(names: &[&str]) -> Self {
Self {
tools: names.iter().map(|name| tool_def(name)).collect(),
dispatched: Mutex::new(Vec::new()),
ops_bound: Mutex::new(false),
mcp_handles_bound: Mutex::new(0),
surface_handles_bound: Mutex::new(0),
}
}
fn dispatched(&self) -> Vec<String> {
self.dispatched.lock().unwrap().clone()
}
}
#[cfg_attr(target_arch = "wasm32", async_trait(?Send))]
#[cfg_attr(not(target_arch = "wasm32"), async_trait)]
impl AgentToolDispatcher for SpyDispatcher {
fn tools(&self) -> Arc<[Arc<ToolDef>]> {
Arc::clone(&self.tools)
}
async fn dispatch(
&self,
call: ToolCallView<'_>,
) -> Result<crate::ops::ToolDispatchOutcome, ToolError> {
self.dispatched.lock().unwrap().push(call.name.to_string());
Ok(crate::ops::ToolDispatchOutcome::from(ToolResult::new(
call.id.to_string(),
"ok".to_string(),
false,
)))
}
fn capabilities(&self) -> DispatcherCapabilities {
DispatcherCapabilities {
ops_lifecycle: true,
}
}
fn bind_ops_lifecycle(
self: Arc<Self>,
_registry: Arc<dyn OpsLifecycleRegistry>,
_owner_bridge_session_id: crate::types::SessionId,
) -> Result<BindOutcome, OpsLifecycleBindError> {
*self.ops_bound.lock().unwrap() = true;
Ok(BindOutcome::Bound(self))
}
fn bind_mcp_server_lifecycle_handle(&self, _handle: Arc<dyn McpServerLifecycleHandle>) {
*self.mcp_handles_bound.lock().unwrap() += 1;
}
fn bind_external_tool_surface_handle(&self, _handle: Arc<dyn ExternalToolSurfaceHandle>) {
*self.surface_handles_bound.lock().unwrap() += 1;
}
}
struct UnsupportedOpsRegistry;
fn unsupported(op: &str) -> OpsLifecycleError {
OpsLifecycleError::Unsupported(op.into())
}
impl OpsLifecycleRegistry for UnsupportedOpsRegistry {
fn register_operation(
&self,
_spec: crate::ops_lifecycle::OperationSpec,
) -> Result<(), OpsLifecycleError> {
Err(unsupported("register_operation"))
}
fn provisioning_succeeded(
&self,
_id: &crate::ops::OperationId,
) -> Result<(), OpsLifecycleError> {
Err(unsupported("provisioning_succeeded"))
}
fn provisioning_failed(
&self,
_id: &crate::ops::OperationId,
_error: String,
) -> Result<(), OpsLifecycleError> {
Err(unsupported("provisioning_failed"))
}
fn peer_ready(
&self,
_id: &crate::ops::OperationId,
_peer: OperationPeerHandle,
) -> Result<(), OpsLifecycleError> {
Err(unsupported("peer_ready"))
}
fn register_watcher(
&self,
_id: &crate::ops::OperationId,
) -> Result<OperationCompletionWatch, OpsLifecycleError> {
Err(unsupported("register_watcher"))
}
fn report_progress(
&self,
_id: &crate::ops::OperationId,
_update: OperationProgressUpdate,
) -> Result<(), OpsLifecycleError> {
Err(unsupported("report_progress"))
}
fn complete_operation(
&self,
_id: &crate::ops::OperationId,
_result: crate::ops::OperationResult,
) -> Result<(), OpsLifecycleError> {
Err(unsupported("complete_operation"))
}
fn fail_operation(
&self,
_id: &crate::ops::OperationId,
_error: String,
) -> Result<(), OpsLifecycleError> {
Err(unsupported("fail_operation"))
}
fn abort_provisioning(
&self,
_id: &crate::ops::OperationId,
_reason: Option<String>,
) -> Result<(), OpsLifecycleError> {
Err(unsupported("abort_provisioning"))
}
fn cancel_operation(
&self,
_id: &crate::ops::OperationId,
_reason: Option<String>,
) -> Result<(), OpsLifecycleError> {
Err(unsupported("cancel_operation"))
}
fn request_retire(&self, _id: &crate::ops::OperationId) -> Result<(), OpsLifecycleError> {
Err(unsupported("request_retire"))
}
fn mark_retired(&self, _id: &crate::ops::OperationId) -> Result<(), OpsLifecycleError> {
Err(unsupported("mark_retired"))
}
fn snapshot(
&self,
_id: &crate::ops::OperationId,
) -> Result<Option<OperationLifecycleSnapshot>, OpsLifecycleError> {
Err(unsupported("snapshot"))
}
fn list_operations(&self) -> Result<Vec<OperationLifecycleSnapshot>, OpsLifecycleError> {
Err(unsupported("list_operations"))
}
fn terminate_owner(&self, _reason: String) -> Result<(), OpsLifecycleError> {
Err(unsupported("terminate_owner"))
}
}
struct NoopMcpLifecycleHandle;
impl McpServerLifecycleHandle for NoopMcpLifecycleHandle {
fn apply_connect_pending(&self, _server_id: &str) -> Result<(), DslTransitionError> {
Ok(())
}
fn apply_connected(&self, _server_id: &str) -> Result<(), DslTransitionError> {
Ok(())
}
fn apply_failed(&self, _server_id: &str, _error: &str) -> Result<(), DslTransitionError> {
Ok(())
}
fn apply_disconnected(&self, _server_id: &str) -> Result<(), DslTransitionError> {
Ok(())
}
fn apply_reload(&self, _server_id: &str) -> Result<(), DslTransitionError> {
Ok(())
}
fn pending_server_ids(&self) -> BTreeSet<String> {
BTreeSet::new()
}
}
struct RejectingSurfaceHandle;
impl RejectingSurfaceHandle {
fn reject(context: &'static str) -> DslTransitionError {
DslTransitionError::guard_rejected(context, "test stub rejects all surface inputs")
}
fn empty_snapshot() -> SurfaceDiagnosticSnapshot {
SurfaceDiagnosticSnapshot {
surface_phase: ExternalToolSurfaceGlobalPhase::Operating,
known_surfaces: BTreeSet::new(),
visible_surfaces: BTreeSet::new(),
snapshot_epoch: 0,
snapshot_aligned_epoch: 0,
has_pending_or_staged: false,
entries: Vec::new(),
}
}
}
impl ExternalToolSurfaceHandle for RejectingSurfaceHandle {
fn apply_surface_input(
&self,
_input: ExternalToolSurfaceInput,
) -> Result<ExternalToolSurfaceTransition, DslTransitionError> {
Err(Self::reject("RejectingSurfaceHandle::apply_surface_input"))
}
fn register(&self, _surface_id: String) -> Result<(), DslTransitionError> {
Err(Self::reject("RejectingSurfaceHandle::register"))
}
fn stage_add(&self, _surface_id: String, _now_ms: u64) -> Result<(), DslTransitionError> {
Err(Self::reject("RejectingSurfaceHandle::stage_add"))
}
fn stage_remove(
&self,
_surface_id: String,
_now_ms: u64,
) -> Result<(), DslTransitionError> {
Err(Self::reject("RejectingSurfaceHandle::stage_remove"))
}
fn stage_reload(
&self,
_surface_id: String,
_now_ms: u64,
) -> Result<(), DslTransitionError> {
Err(Self::reject("RejectingSurfaceHandle::stage_reload"))
}
fn apply_boundary(
&self,
_surface_id: String,
_now_ms: u64,
_staged_intent_sequence: u64,
_applied_at_turn: u64,
) -> Result<(), DslTransitionError> {
Err(Self::reject("RejectingSurfaceHandle::apply_boundary"))
}
fn mark_pending_succeeded(
&self,
_surface_id: String,
_pending_task_sequence: u64,
_staged_intent_sequence: u64,
) -> Result<(), DslTransitionError> {
Err(Self::reject(
"RejectingSurfaceHandle::mark_pending_succeeded",
))
}
fn mark_pending_failed(
&self,
_surface_id: String,
_pending_task_sequence: u64,
_staged_intent_sequence: u64,
_cause: crate::tool_scope::ExternalToolSurfaceFailureCause,
) -> Result<(), DslTransitionError> {
Err(Self::reject("RejectingSurfaceHandle::mark_pending_failed"))
}
fn call_started(&self, _surface_id: String) -> Result<(), DslTransitionError> {
Err(Self::reject("RejectingSurfaceHandle::call_started"))
}
fn call_finished(&self, _surface_id: String) -> Result<(), DslTransitionError> {
Err(Self::reject("RejectingSurfaceHandle::call_finished"))
}
fn finalize_removal_clean(&self, _surface_id: String) -> Result<(), DslTransitionError> {
Err(Self::reject(
"RejectingSurfaceHandle::finalize_removal_clean",
))
}
fn finalize_removal_forced(&self, _surface_id: String) -> Result<(), DslTransitionError> {
Err(Self::reject(
"RejectingSurfaceHandle::finalize_removal_forced",
))
}
fn snapshot_aligned(&self, _epoch: u64) -> Result<(), DslTransitionError> {
Err(Self::reject("RejectingSurfaceHandle::snapshot_aligned"))
}
fn shutdown_surface(&self) -> Result<(), DslTransitionError> {
Err(Self::reject("RejectingSurfaceHandle::shutdown_surface"))
}
fn surface_snapshot(&self, _surface_id: &str) -> Option<SurfaceSnapshot> {
None
}
fn diagnostic_snapshot(&self) -> SurfaceDiagnosticSnapshot {
Self::empty_snapshot()
}
fn visible_surfaces(&self) -> BTreeSet<String> {
BTreeSet::new()
}
fn removing_surfaces(&self) -> BTreeSet<String> {
BTreeSet::new()
}
fn pending_surfaces(&self) -> BTreeSet<String> {
BTreeSet::new()
}
fn has_pending_or_staged(&self) -> bool {
false
}
fn snapshot_epoch(&self) -> u64 {
0
}
fn snapshot_aligned_epoch(&self) -> u64 {
0
}
}
fn allow_list(names: &[&str]) -> ToolExecutionPolicy {
ToolExecutionPolicy::resolve(ToolAccessPolicy::AllowList(names.iter().copied().collect()))
.expect("allow list resolves")
}
fn deny_list(names: &[&str]) -> ToolExecutionPolicy {
ToolExecutionPolicy::resolve(ToolAccessPolicy::DenyList(names.iter().copied().collect()))
.expect("deny list resolves")
}
fn read_only() -> ToolExecutionPolicy {
ToolExecutionPolicy::resolve(ToolAccessPolicy::ReadOnly).expect("read-only resolves")
}
struct DeclaringDispatcher {
inner: Arc<SpyDispatcher>,
classes: Vec<(String, ToolMutationClass)>,
}
impl DeclaringDispatcher {
fn new(declared: &[(&str, ToolMutationClass)], undeclared: &[&str]) -> Self {
let names: Vec<&str> = declared
.iter()
.map(|(name, _)| *name)
.chain(undeclared.iter().copied())
.collect();
Self {
inner: Arc::new(SpyDispatcher::new(&names)),
classes: declared
.iter()
.map(|(name, class)| ((*name).to_string(), *class))
.collect(),
}
}
fn dispatched(&self) -> Vec<String> {
self.inner.dispatched()
}
}
#[cfg_attr(target_arch = "wasm32", async_trait(?Send))]
#[cfg_attr(not(target_arch = "wasm32"), async_trait)]
impl AgentToolDispatcher for DeclaringDispatcher {
fn tools(&self) -> Arc<[Arc<ToolDef>]> {
self.inner.tools()
}
fn tool_mutation_class(&self, tool_name: &str) -> ToolMutationClass {
self.classes
.iter()
.find(|(name, _)| name == tool_name)
.map(|(_, class)| *class)
.unwrap_or_default()
}
async fn dispatch(
&self,
call: ToolCallView<'_>,
) -> Result<crate::ops::ToolDispatchOutcome, ToolError> {
self.inner.dispatch(call).await
}
async fn dispatch_with_context(
&self,
call: ToolCallView<'_>,
context: &ToolDispatchContext,
) -> Result<crate::ops::ToolDispatchOutcome, ToolError> {
self.inner.dispatch_with_context(call, context).await
}
}
async fn dispatch_named<T: AgentToolDispatcher + ?Sized>(
dispatcher: &T,
name: &str,
) -> Result<crate::ops::ToolDispatchOutcome, ToolError> {
let args = empty_args();
let call = ToolCallView {
id: "call-1",
name,
args: &args,
};
dispatcher.dispatch(call).await
}
async fn dispatch_named_with_context<T: AgentToolDispatcher + ?Sized>(
dispatcher: &T,
name: &str,
) -> Result<crate::ops::ToolDispatchOutcome, ToolError> {
let args = empty_args();
let call = ToolCallView {
id: "call-1",
name,
args: &args,
};
dispatcher
.dispatch_with_context(call, &ToolDispatchContext::default())
.await
}
#[test]
fn resolve_inherit_fails_closed_with_typed_error() {
let err = ToolExecutionPolicy::resolve(ToolAccessPolicy::Inherit)
.expect_err("inherit must not resolve at the execution seam");
assert_eq!(err, ToolExecutionPolicyError::UnresolvedInherit);
assert_eq!(err.error_code(), "tool_execution_policy_unresolved_inherit");
}
#[test]
fn resolve_allow_and_deny_lists_carry_over() {
assert!(allow_list(&["a"]).permits("a"));
assert!(!allow_list(&["a"]).permits("b"));
assert!(!deny_list(&["a"]).permits("a"));
assert!(deny_list(&["a"]).permits("b"));
assert!(ToolExecutionPolicy::unrestricted().permits("anything"));
assert!(ToolExecutionPolicy::unrestricted().is_unrestricted());
assert!(!allow_list(&["a"]).is_unrestricted());
assert!(!deny_list(&["a"]).is_unrestricted());
}
#[test]
fn gated_dispatcher_preserves_tools_and_catalog_byte_identically() {
let inner = Arc::new(SpyDispatcher::new(&["alpha", "beta", "gamma"]));
let inner_tools = inner.tools();
let inner_catalog = inner.tool_catalog();
let gated = ExecutionPolicyGatedDispatcher::new(Arc::clone(&inner), allow_list(&["alpha"]));
let gated_tools = gated.tools();
assert_eq!(gated_tools.len(), inner_tools.len());
for (gated_tool, inner_tool) in gated_tools.iter().zip(inner_tools.iter()) {
assert!(Arc::ptr_eq(gated_tool, inner_tool));
}
let gated_catalog = gated.tool_catalog();
assert_eq!(gated_catalog.len(), inner_catalog.len());
for (gated_entry, inner_entry) in gated_catalog.iter().zip(inner_catalog.iter()) {
assert!(Arc::ptr_eq(&gated_entry.tool, &inner_entry.tool));
}
assert_eq!(
gated.tool_catalog_capabilities(),
inner.tool_catalog_capabilities()
);
assert_eq!(gated.capabilities(), inner.capabilities());
}
#[tokio::test]
async fn gated_dispatcher_denies_plan_resolution_and_resolved_dispatch() {
let inner = Arc::new(SpyDispatcher::new(&["alpha", "beta"]));
let gated = ExecutionPolicyGatedDispatcher::new(Arc::clone(&inner), allow_list(&["alpha"]));
let args = empty_args();
let call = ToolCallView {
id: "call-plan",
name: "beta",
args: &args,
};
let resolution = crate::ToolExecutionResolutionContext::new(
crate::ToolDeadlineChain::new(vec![crate::ToolDeadlineContributor::finite(
crate::ToolDeadlineOwner::CoreToolDispatch,
std::time::Duration::from_secs(600),
)])
.unwrap(),
);
let error = gated
.resolve_execution_plan(call, &ToolDispatchContext::default(), &resolution)
.expect_err("policy-denied tools must not resolve execution preparation");
assert_eq!(
error,
crate::ToolExecutionResolutionError::AccessDenied {
tool_name: "beta".to_string(),
}
);
let bypass_plan = crate::ToolExecutionContract::default()
.resolve_default(resolution.deadlines().clone())
.expect("test fast plan resolves");
let dispatch_error = gated
.dispatch_resolved_with_context(call, &ToolDispatchContext::default(), &bypass_plan)
.await
.expect_err("policy gate must be rechecked at resolved dispatch");
assert!(matches!(dispatch_error, ToolError::AccessDenied { .. }));
assert!(inner.dispatched().is_empty());
}
#[tokio::test]
async fn allow_list_matrix_permits_listed_denies_rest() {
let inner = Arc::new(SpyDispatcher::new(&["alpha", "beta"]));
let gated = ExecutionPolicyGatedDispatcher::new(Arc::clone(&inner), allow_list(&["alpha"]));
let outcome = dispatch_named(&gated, "alpha")
.await
.expect("allow-listed tool must dispatch");
assert!(!outcome.result.is_error);
let err = dispatch_named(&gated, "beta")
.await
.expect_err("non-listed known tool must be denied");
assert_eq!(err, ToolError::access_denied("beta"));
assert_eq!(err.error_code(), "access_denied");
let err = dispatch_named(&gated, "missing")
.await
.expect_err("unknown tool must not dispatch");
assert_eq!(err, ToolError::not_found("missing"));
assert_eq!(inner.dispatched(), vec!["alpha".to_string()]);
}
#[tokio::test]
async fn deny_list_matrix_denies_listed_permits_rest() {
let inner = Arc::new(SpyDispatcher::new(&["alpha", "beta"]));
let gated = ExecutionPolicyGatedDispatcher::new(Arc::clone(&inner), deny_list(&["beta"]));
dispatch_named_with_context(&gated, "alpha")
.await
.expect("non-denied tool must dispatch");
let err = dispatch_named_with_context(&gated, "beta")
.await
.expect_err("deny-listed tool must be denied");
assert_eq!(err, ToolError::access_denied("beta"));
let gated_ghost =
ExecutionPolicyGatedDispatcher::new(Arc::clone(&inner), deny_list(&["ghost"]));
let err = dispatch_named(&gated_ghost, "ghost")
.await
.expect_err("unknown deny-listed tool must not dispatch");
assert_eq!(err, ToolError::not_found("ghost"));
dispatch_named(&gated, "gamma")
.await
.expect("policy-permitted unknown name forwards to inner");
assert_eq!(
inner.dispatched(),
vec!["alpha".to_string(), "gamma".to_string()]
);
}
#[tokio::test]
async fn memory_search_denied_when_absent_from_allow_list() {
let inner = Arc::new(SpyDispatcher::new(&["memory_search", "read_file"]));
let gated =
ExecutionPolicyGatedDispatcher::new(Arc::clone(&inner), allow_list(&["read_file"]));
let err = dispatch_named(&gated, "memory_search")
.await
.expect_err("memory_search absent from allow list must be denied");
assert_eq!(err, ToolError::access_denied("memory_search"));
assert!(inner.dispatched().is_empty());
}
#[tokio::test]
async fn read_only_admits_declared_reads_and_refuses_everything_else() {
let inner = Arc::new(DeclaringDispatcher::new(
&[
("datetime", ToolMutationClass::ReadOnly),
("shell", ToolMutationClass::Mutating),
],
&["mcp_unknown_tool"],
));
let gated = ExecutionPolicyGatedDispatcher::new(Arc::clone(&inner), read_only());
let outcome = dispatch_named(&gated, "datetime")
.await
.expect("declared read-only tool must dispatch");
assert!(!outcome.result.is_error);
let err = dispatch_named(&gated, "shell")
.await
.expect_err("declared mutating tool must be denied");
assert_eq!(err, ToolError::access_denied("shell"));
assert_eq!(err.error_code(), "access_denied");
let err = dispatch_named(&gated, "mcp_unknown_tool")
.await
.expect_err("undeclared tool must be denied under read-only intent");
assert_eq!(err, ToolError::access_denied("mcp_unknown_tool"));
assert_eq!(inner.dispatched(), vec!["datetime".to_string()]);
}
#[tokio::test]
async fn read_only_refusal_precedes_execution_on_every_dispatch_entry_point() {
let inner = Arc::new(DeclaringDispatcher::new(
&[("shell", ToolMutationClass::Mutating)],
&[],
));
let gated = ExecutionPolicyGatedDispatcher::new(Arc::clone(&inner), read_only());
dispatch_named(&gated, "shell")
.await
.expect_err("dispatch must deny");
dispatch_named_with_context(&gated, "shell")
.await
.expect_err("dispatch_with_context must deny");
let args = empty_args();
let call = ToolCallView {
id: "call-1",
name: "shell",
args: &args,
};
let resolution = crate::ToolExecutionResolutionContext::new(
crate::ToolDeadlineChain::new(vec![crate::ToolDeadlineContributor::finite(
crate::ToolDeadlineOwner::CoreToolDispatch,
std::time::Duration::from_secs(600),
)])
.expect("deadline chain"),
);
gated
.resolve_execution_plan(call, &ToolDispatchContext::default(), &resolution)
.expect_err("plan resolution must deny before any execution plan exists");
assert!(
inner.dispatched().is_empty(),
"no read-only denial may reach the inner dispatcher"
);
}
#[test]
fn read_only_policy_denies_when_no_declaration_is_supplied() {
let policy = read_only();
assert!(!policy.permits("datetime"));
assert!(policy.permits_call("datetime", ToolMutationClass::ReadOnly));
assert!(!policy.permits_call("datetime", ToolMutationClass::Mutating));
assert!(!policy.permits_call("datetime", ToolMutationClass::Unknown));
assert!(policy.is_read_only_intent());
assert!(!policy.is_unrestricted());
}
#[test]
fn name_list_policies_ignore_mutation_declarations() {
assert!(allow_list(&["alpha"]).permits_call("alpha", ToolMutationClass::Mutating));
assert!(!allow_list(&["alpha"]).permits_call("beta", ToolMutationClass::ReadOnly));
assert!(deny_list(&["beta"]).permits_call("alpha", ToolMutationClass::Unknown));
assert!(!deny_list(&["beta"]).permits_call("beta", ToolMutationClass::ReadOnly));
assert!(
ToolExecutionPolicy::unrestricted().permits_call("beta", ToolMutationClass::Unknown)
);
}
#[test]
fn conjunctive_constraints_never_widen_each_other() {
let policy = ToolAccessPolicy::AllowList(["a"].into_iter().collect())
.conjoin(ToolAccessPolicy::ReadOnly)
.expect("concrete policies conjoin");
let resolved = ToolExecutionPolicy::resolve(policy).expect("constraints resolve");
assert!(resolved.permits_call("a", ToolMutationClass::ReadOnly));
assert!(!resolved.permits_call("a", ToolMutationClass::Mutating));
assert!(!resolved.permits_call("b", ToolMutationClass::ReadOnly));
let denied = ToolAccessPolicy::AllowList(["a"].into_iter().collect())
.conjoin(ToolAccessPolicy::DenyList(["a"].into_iter().collect()))
.expect("concrete policies conjoin");
assert!(
!ToolExecutionPolicy::resolve(denied)
.expect("constraints resolve")
.permits_call("a", ToolMutationClass::ReadOnly)
);
}
#[test]
fn empty_constraint_set_is_rejected() {
assert_eq!(
ToolExecutionPolicy::resolve(ToolAccessPolicy::Constraints(Vec::new()))
.expect_err("empty constraints are not unrestricted"),
ToolExecutionPolicyError::EmptyConstraints
);
}
#[tokio::test]
async fn read_only_gate_keeps_the_llm_visible_tool_list_unchanged() {
let inner = Arc::new(DeclaringDispatcher::new(
&[
("datetime", ToolMutationClass::ReadOnly),
("shell", ToolMutationClass::Mutating),
],
&[],
));
let gated = ExecutionPolicyGatedDispatcher::new(Arc::clone(&inner), read_only());
let before: Vec<String> = inner
.tools()
.iter()
.map(|tool| tool.name.to_string())
.collect();
let after: Vec<String> = gated
.tools()
.iter()
.map(|tool| tool.name.to_string())
.collect();
assert_eq!(
before, after,
"read-only gating must not change the prompt-cache prefix"
);
}
#[tokio::test]
async fn bind_ops_lifecycle_rewrap_keeps_gate_and_registry_binding() {
let inner = Arc::new(SpyDispatcher::new(&["alpha", "beta"]));
let gated: Arc<ExecutionPolicyGatedDispatcher<SpyDispatcher>> = Arc::new(
ExecutionPolicyGatedDispatcher::new(Arc::clone(&inner), allow_list(&["alpha"])),
);
let inner_probe = Arc::downgrade(&inner);
drop(inner);
let outcome = gated
.bind_ops_lifecycle(
Arc::new(UnsupportedOpsRegistry),
crate::types::SessionId::new(),
)
.expect("bind must succeed through the gate");
assert!(outcome.was_bound(), "inner binding must be applied");
let rebound = outcome.into_dispatcher();
let inner_alive = inner_probe
.upgrade()
.expect("inner dispatcher must survive rebind");
assert!(
*inner_alive.ops_bound.lock().unwrap(),
"ops registry binding must reach the inner dispatcher"
);
let err = dispatch_named_with_context(rebound.as_ref(), "beta")
.await
.expect_err("gate must survive bind_ops_lifecycle re-wrap");
assert_eq!(err, ToolError::access_denied("beta"));
dispatch_named_with_context(rebound.as_ref(), "alpha")
.await
.expect("allow-listed tool must still dispatch after re-wrap");
}
#[test]
fn bind_ops_lifecycle_shared_wrapper_reports_shared_ownership() {
let inner = Arc::new(SpyDispatcher::new(&["alpha"]));
let gated: Arc<ExecutionPolicyGatedDispatcher<SpyDispatcher>> = Arc::new(
ExecutionPolicyGatedDispatcher::new(Arc::clone(&inner), allow_list(&["alpha"])),
);
let extra_handle = Arc::clone(&gated);
let err = match gated.bind_ops_lifecycle(
Arc::new(UnsupportedOpsRegistry),
crate::types::SessionId::new(),
) {
Ok(_) => panic!("shared wrapper ownership must refuse rebind"),
Err(err) => err,
};
assert_eq!(err, OpsLifecycleBindError::SharedOwnership);
drop(extra_handle);
}
#[test]
fn both_handle_binds_reach_inner_dispatcher() {
let inner = Arc::new(SpyDispatcher::new(&["alpha"]));
let gated = ExecutionPolicyGatedDispatcher::new(Arc::clone(&inner), allow_list(&["alpha"]));
gated.bind_mcp_server_lifecycle_handle(Arc::new(NoopMcpLifecycleHandle));
gated.bind_external_tool_surface_handle(Arc::new(RejectingSurfaceHandle));
assert_eq!(
*inner.mcp_handles_bound.lock().unwrap(),
1,
"bind_mcp_server_lifecycle_handle must forward to inner"
);
assert_eq!(
*inner.surface_handles_bound.lock().unwrap(),
1,
"bind_external_tool_surface_handle must forward to inner"
);
}
struct FixedConsequenceSnapshot {
verdict: ToolConsequenceVerdict,
}
impl ToolConsequencePolicySnapshot for FixedConsequenceSnapshot {
fn provenance(&self) -> PolicyEvaluationProvenance {
PolicyEvaluationProvenance {
revision: PolicyRevision(7),
digest: PolicyDigest::from_canonical_bytes(b"fixed-test-policy"),
}
}
fn evaluate(&self, _request: &ToolConsequenceRequest) -> ToolConsequenceVerdict {
self.verdict.clone()
}
}
struct FixedConsequenceProvider {
provider_id: PolicyProviderId,
snapshot: Arc<dyn ToolConsequencePolicySnapshot>,
}
impl ToolConsequenceNarrowingPolicy for FixedConsequenceProvider {
fn provider_id(&self) -> &PolicyProviderId {
&self.provider_id
}
fn generation(&self) -> PolicyProviderGeneration {
PolicyProviderGeneration(1)
}
fn snapshot(
&self,
_policy_id: &PolicyId,
) -> Result<Arc<dyn ToolConsequencePolicySnapshot>, ToolConsequenceFailure> {
Ok(Arc::clone(&self.snapshot))
}
}
fn bound_consequence_policy(
verdict: ToolConsequenceVerdict,
deadline: std::time::Duration,
) -> BoundToolConsequencePolicy {
let provider_id = PolicyProviderId::new("test-provider").expect("provider id");
let policy_id = PolicyId::new("test-policy").expect("policy id");
let provider: Arc<dyn ToolConsequenceNarrowingPolicy> =
Arc::new(FixedConsequenceProvider {
provider_id: provider_id.clone(),
snapshot: Arc::new(FixedConsequenceSnapshot { verdict }),
});
let registry = Arc::new(
ToolConsequencePolicyRegistry::new(
vec![provider],
PolicyEvaluationSupervisorConfig {
workers_per_provider: 1,
queue_capacity_per_provider: 1,
evaluation_deadline: deadline,
},
None,
)
.expect("registry"),
);
registry
.bind(
MobMemberBinding {
mob_id: "mob".to_string(),
role: "worker".to_string(),
member: "member".to_string(),
},
provider_id,
policy_id,
)
.expect("binding")
}
#[tokio::test]
async fn application_denial_is_narrow_only_and_never_enters_inner_dispatcher() {
let inner = Arc::new(SpyDispatcher::new(&["alpha"]));
let gated = ExecutionPolicyGatedDispatcher::new(
Arc::clone(&inner),
ToolExecutionPolicy::unrestricted(),
)
.with_consequence_policy(bound_consequence_policy(
ToolConsequenceVerdict::Deny(ToolConsequenceDenial::new(
"test_denied",
"denied by test policy",
)),
std::time::Duration::from_millis(100),
));
let error = dispatch_named_with_context(&gated, "alpha")
.await
.expect_err("application policy must deny");
assert!(matches!(error, ToolError::PolicyDenied { .. }));
assert!(inner.dispatched().is_empty());
}
#[tokio::test]
async fn static_denial_precedes_application_policy() {
let inner = Arc::new(SpyDispatcher::new(&["alpha", "beta"]));
let gated = ExecutionPolicyGatedDispatcher::new(Arc::clone(&inner), allow_list(&["alpha"]))
.with_consequence_policy(bound_consequence_policy(
ToolConsequenceVerdict::Allow,
std::time::Duration::from_millis(100),
));
let error = dispatch_named_with_context(&gated, "beta")
.await
.expect_err("static policy must remain authoritative");
assert_eq!(error, ToolError::access_denied("beta"));
assert!(inner.dispatched().is_empty());
}
#[tokio::test]
async fn dispatch_admission_is_outermost_and_still_intersects_ordinary_policy() {
let inner = Arc::new(SpyDispatcher::new(&["alpha", "beta"]));
let admission = Arc::new(BlockingAdmission::new());
let gated = Arc::new(
ExecutionPolicyGatedDispatcher::new(Arc::clone(&inner), allow_list(&["alpha"]))
.with_dispatch_admission(Arc::clone(&admission) as Arc<dyn ToolDispatchAdmission>),
);
let pending = tokio::spawn({
let gated = Arc::clone(&gated);
async move { dispatch_named_with_context(gated.as_ref(), "alpha").await }
});
tokio::task::yield_now().await;
assert!(
inner.dispatched().is_empty(),
"inner dispatcher must not observe a provisional tool call"
);
let denied = tokio::spawn({
let gated = Arc::clone(&gated);
async move { dispatch_named_with_context(gated.as_ref(), "beta").await }
});
tokio::task::yield_now().await;
assert!(inner.dispatched().is_empty());
admission.release();
pending
.await
.expect("dispatch task")
.expect("released exact operation must dispatch");
let denied = denied
.await
.expect("denied dispatch task")
.expect_err("ordinary policy must still narrow released admission");
assert_eq!(denied, ToolError::access_denied("beta"));
assert_eq!(inner.dispatched(), vec!["alpha"]);
let outcomes = admission.outcomes.lock().unwrap().clone();
assert_eq!(outcomes.len(), 2);
assert!(outcomes.contains(&crate::LiveBridgeEffectOutcome::Committed));
assert!(outcomes.contains(&crate::LiveBridgeEffectOutcome::Failed));
}
#[tokio::test]
async fn entered_inner_dispatch_error_is_unknown_not_retriable_failure() {
struct FailingDispatcher(SpyDispatcher);
#[async_trait::async_trait]
impl AgentToolDispatcher for FailingDispatcher {
fn tools(&self) -> Arc<[Arc<ToolDef>]> {
self.0.tools()
}
async fn dispatch(
&self,
_call: ToolCallView<'_>,
) -> Result<crate::ops::ToolDispatchOutcome, ToolError> {
Err(ToolError::ExecutionFailed {
message: "physical outcome unavailable".to_string(),
})
}
}
let admission = Arc::new(BlockingAdmission::new());
admission.release();
let gated = ExecutionPolicyGatedDispatcher::new(
Arc::new(FailingDispatcher(SpyDispatcher::new(&["external"]))),
ToolExecutionPolicy::unrestricted(),
)
.with_dispatch_admission(Arc::clone(&admission) as Arc<dyn ToolDispatchAdmission>);
dispatch_named(&gated, "external")
.await
.expect_err("inner dispatcher reports an error");
assert_eq!(
*admission.outcomes.lock().unwrap(),
vec![crate::LiveBridgeEffectOutcome::Unknown],
"an error after entering physical dispatch cannot be relabeled as a retryable failure"
);
}
#[tokio::test]
async fn operation_context_admission_receives_owner_declared_effect_kind() {
struct MemoryEffectDispatcher(SpyDispatcher);
#[async_trait::async_trait]
impl AgentToolDispatcher for MemoryEffectDispatcher {
fn tools(&self) -> Arc<[Arc<ToolDef>]> {
self.0.tools()
}
fn live_bridge_effect_kind(&self, _tool_name: &str) -> crate::LiveBridgeEffectKind {
crate::LiveBridgeEffectKind::ReadOnlyMemorySnapshot
}
async fn dispatch(
&self,
call: ToolCallView<'_>,
) -> Result<crate::ops::ToolDispatchOutcome, ToolError> {
self.0.dispatch(call).await
}
}
let inner = Arc::new(MemoryEffectDispatcher(SpyDispatcher::new(&[
"memory_search",
])));
let admission = Arc::new(BlockingAdmission::new());
admission.release();
let gated = ExecutionPolicyGatedDispatcher::new(inner, ToolExecutionPolicy::unrestricted());
let context = ToolDispatchContext::default().with_live_bridge_admission(
crate::LiveBridgeToolDispatchAdmission::new(
"operation-1",
Arc::clone(&admission) as Arc<dyn ToolDispatchAdmission>,
),
);
let args = empty_args();
gated
.dispatch_with_context(
ToolCallView {
id: "call-1",
name: "memory_search",
args: &args,
},
&context,
)
.await
.expect("owner-classified dispatch");
assert_eq!(
*admission.effect_kinds.lock().unwrap(),
vec![crate::LiveBridgeEffectKind::ReadOnlyMemorySnapshot]
);
assert_eq!(
*admission.outcomes.lock().unwrap(),
vec![crate::LiveBridgeEffectOutcome::Committed]
);
}
struct WedgedConsequenceSnapshot;
impl ToolConsequencePolicySnapshot for WedgedConsequenceSnapshot {
fn provenance(&self) -> PolicyEvaluationProvenance {
PolicyEvaluationProvenance {
revision: PolicyRevision(1),
digest: PolicyDigest::from_canonical_bytes(b"wedged"),
}
}
fn evaluate(&self, _request: &ToolConsequenceRequest) -> ToolConsequenceVerdict {
std::thread::sleep(std::time::Duration::from_secs(1));
ToolConsequenceVerdict::Allow
}
}
#[tokio::test]
async fn wedged_evaluator_deadlines_and_partition_then_fails_fast() {
let provider_id = PolicyProviderId::new("wedged-provider").expect("provider id");
let policy_id = PolicyId::new("policy").expect("policy id");
let provider: Arc<dyn ToolConsequenceNarrowingPolicy> =
Arc::new(FixedConsequenceProvider {
provider_id: provider_id.clone(),
snapshot: Arc::new(WedgedConsequenceSnapshot),
});
let registry = Arc::new(
ToolConsequencePolicyRegistry::new(
vec![provider],
PolicyEvaluationSupervisorConfig {
workers_per_provider: 1,
queue_capacity_per_provider: 1,
evaluation_deadline: std::time::Duration::from_millis(5),
},
None,
)
.expect("registry"),
);
let policy = registry
.bind(
MobMemberBinding {
mob_id: "mob".to_string(),
role: "worker".to_string(),
member: "member".to_string(),
},
provider_id,
policy_id,
)
.expect("binding");
let inner = Arc::new(SpyDispatcher::new(&["alpha"]));
let gated = ExecutionPolicyGatedDispatcher::new(
Arc::clone(&inner),
ToolExecutionPolicy::unrestricted(),
)
.with_consequence_policy(policy);
let first = dispatch_named_with_context(&gated, "alpha")
.await
.expect_err("wedged policy must deadline");
assert!(matches!(
first,
ToolError::PolicyIndeterminate {
failure: ToolConsequenceFailure::DeadlineExceeded { .. }
}
));
let second = dispatch_named_with_context(&gated, "alpha")
.await
.expect_err("unhealthy partition must fail fast");
assert!(matches!(
second,
ToolError::PolicyIndeterminate {
failure: ToolConsequenceFailure::MechanicallyUnhealthy { .. }
}
));
assert!(inner.dispatched().is_empty());
}
}