1use gate4agent_tool_protocol::*;
2use gate4agent_types::{AgentInstanceId, SessionGeneration};
3use std::collections::{BTreeMap, VecDeque};
4use std::fmt;
5
6#[derive(Clone, Eq, PartialEq)]
7struct RequestState {
8 request: AcceptedRequest,
9 snapshot: CapabilityRequestSnapshot,
10}
11
12#[derive(Clone, Eq, PartialEq)]
13struct AcceptedRequest {
14 key: CapabilityRequestKey,
15 instance_id: AgentInstanceId,
16 generation: SessionGeneration,
17 provider_id: ToolProviderId,
18 capability_id: ToolCapabilityId,
19 resource_scope_id: ResourceScopeId,
20 approval_summary: String,
21 deadline_tick: u64,
22 payload: Vec<u8>,
23}
24
25impl From<ConsumerBoundCapabilityRequest> for AcceptedRequest {
26 fn from(envelope: ConsumerBoundCapabilityRequest) -> Self {
27 Self {
28 key: envelope.key(),
29 instance_id: envelope.request.instance_id,
30 generation: envelope.request.generation,
31 provider_id: envelope.request.provider_id,
32 capability_id: envelope.request.capability_id,
33 resource_scope_id: envelope.request.resource_scope_id,
34 approval_summary: envelope.request.approval_summary,
35 deadline_tick: envelope.request.deadline_tick,
36 payload: envelope.request.payload,
37 }
38 }
39}
40
41#[derive(Clone, Copy)]
42enum RequestCloseKind {
43 Instance,
44 Client,
45}
46
47#[derive(Clone, Debug, Eq, PartialEq)]
48pub enum ToolEngineError {
49 Validation(ToolValidationError),
50 DuplicateProvider {
51 provider_id: ToolProviderId,
52 },
53 UnknownRuntimeProvider {
54 provider_id: ToolProviderId,
55 },
56 DuplicateRequest {
57 request_key: CapabilityRequestKey,
58 },
59 ProviderCapacityExceeded,
60 PolicyCapacityExceeded,
61 RequestCapacityExceeded,
62 ClientRequestCapacityExceeded {
63 consumer_id: ConsumerId,
64 actor_id: ToolActorId,
65 max: usize,
66 },
67 EffectCapacityExceeded,
68 EffectSequenceExhausted,
69 UnknownPolicyProvider {
70 provider_id: ToolProviderId,
71 },
72 UnknownPolicyCapability {
73 provider_id: ToolProviderId,
74 capability_id: ToolCapabilityId,
75 },
76 ProviderOwnerMismatch {
77 provider_id: ToolProviderId,
78 owner: ConsumerId,
79 requested_consumer: ConsumerId,
80 },
81 UnknownPolicyInstance {
82 instance_id: AgentInstanceId,
83 },
84 InactivePolicyInstance {
85 instance_id: AgentInstanceId,
86 },
87 PolicyGenerationMismatch {
88 instance_id: AgentInstanceId,
89 current: SessionGeneration,
90 requested: SessionGeneration,
91 },
92 UnknownRequest {
93 request_key: CapabilityRequestKey,
94 },
95 RequestNotAwaitingApproval {
96 request_key: CapabilityRequestKey,
97 },
98 ApprovalScopeMismatch {
99 request_key: CapabilityRequestKey,
100 },
101 ApprovalNonceMismatch {
102 request_key: CapabilityRequestKey,
103 expected: u64,
104 actual: u64,
105 },
106 ApprovalGenerationStale {
107 current: SessionGeneration,
108 actual: SessionGeneration,
109 },
110 ApprovalDeadlineElapsed {
111 request_key: CapabilityRequestKey,
112 },
113 ClockRegressed {
114 current_tick: u64,
115 requested_tick: u64,
116 },
117 GenerationRegressed {
118 instance_id: AgentInstanceId,
119 current: SessionGeneration,
120 requested: SessionGeneration,
121 },
122 AuthoritySequenceRegressed {
123 current: u64,
124 requested: u64,
125 },
126 CounterExhausted {
127 counter: &'static str,
128 },
129}
130
131impl From<ToolValidationError> for ToolEngineError {
132 fn from(error: ToolValidationError) -> Self {
133 Self::Validation(error)
134 }
135}
136
137impl fmt::Display for ToolEngineError {
138 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
139 write!(formatter, "tool engine rejected transition: {self:?}")
140 }
141}
142
143impl std::error::Error for ToolEngineError {}
144
145#[derive(Clone, Eq, PartialEq)]
150pub struct ToolEngine {
151 revision: u64,
152 current_tick: u64,
153 generations: BTreeMap<AgentInstanceId, SessionGeneration>,
154 instance_states: BTreeMap<AgentInstanceId, ToolInstanceState>,
155 providers: BTreeMap<ToolProviderId, CapabilityProviderDescriptor>,
156 grants: BTreeMap<PolicyKey, GrantMode>,
157 requests: BTreeMap<CapabilityRequestKey, RequestState>,
158 effects: Vec<CapabilityEffectEnvelope>,
159 completions: Vec<CapabilityCompletionEnvelope>,
160 audit_events: VecDeque<ToolAuditEvent>,
161 dropped_audit_events: u64,
162 revision_overflow_count: u64,
163 dropped_completions: u64,
164 dropped_completions_since_drain: u64,
165 last_authority_sequence: u64,
166 effect_sequence_exhausted: bool,
167 completion_sequence_exhausted: bool,
168 audit_sequence_exhausted: bool,
169 next_request_sequence: u64,
170 next_operation_id: u64,
171 next_effect_sequence: u64,
172 next_completion_sequence: u64,
173 next_audit_sequence: u64,
174}
175
176impl ToolEngine {
177 pub fn new() -> Self {
178 Self {
179 revision: 0,
180 current_tick: 0,
181 generations: BTreeMap::new(),
182 instance_states: BTreeMap::new(),
183 providers: BTreeMap::new(),
184 grants: BTreeMap::new(),
185 requests: BTreeMap::new(),
186 effects: Vec::new(),
187 completions: Vec::new(),
188 audit_events: VecDeque::new(),
189 dropped_audit_events: 0,
190 revision_overflow_count: 0,
191 dropped_completions: 0,
192 dropped_completions_since_drain: 0,
193 last_authority_sequence: 0,
194 effect_sequence_exhausted: false,
195 completion_sequence_exhausted: false,
196 audit_sequence_exhausted: false,
197 next_request_sequence: 1,
198 next_operation_id: 1,
199 next_effect_sequence: 1,
200 next_completion_sequence: 1,
201 next_audit_sequence: 1,
202 }
203 }
204
205 pub fn register_provider(
206 &mut self,
207 mut descriptor: CapabilityProviderDescriptor,
208 ) -> Result<(), ToolEngineError> {
209 descriptor.validate()?;
210 if self.providers.contains_key(&descriptor.id) {
211 return Err(ToolEngineError::DuplicateProvider {
212 provider_id: descriptor.id,
213 });
214 }
215 if self.providers.len() >= TOOL_PROVIDERS_MAX {
216 return Err(ToolEngineError::ProviderCapacityExceeded);
217 }
218 descriptor
219 .capabilities
220 .sort_by(|left, right| left.id.cmp(&right.id));
221 let provider_id = descriptor.id.clone();
222 let owner = descriptor.owner.clone();
223 let capability_count = descriptor.capabilities.len();
224 self.providers.insert(provider_id.clone(), descriptor);
225 self.bump_revision();
226 self.emit_audit(
227 None,
228 ToolAuditEventKind::ProviderRegistered {
229 provider_id,
230 owner,
231 capability_count,
232 },
233 );
234 Ok(())
235 }
236
237 pub fn provider_exists(&self, provider_id: &ToolProviderId) -> bool {
238 self.providers.contains_key(provider_id)
239 }
240
241 pub fn provider_descriptor(
242 &self,
243 provider_id: &ToolProviderId,
244 ) -> Option<&CapabilityProviderDescriptor> {
245 self.providers.get(provider_id)
246 }
247
248 pub fn detach_provider_runtime(
249 &mut self,
250 provider_id: &ToolProviderId,
251 ) -> Result<usize, ToolEngineError> {
252 if !self.provider_exists(provider_id) {
253 return Err(ToolEngineError::UnknownRuntimeProvider {
254 provider_id: provider_id.clone(),
255 });
256 }
257 let targets = self
258 .requests
259 .iter()
260 .filter_map(|(request_key, state)| {
261 (&state.request.provider_id == provider_id && !state.snapshot.status.is_terminal())
262 .then_some(request_key.clone())
263 })
264 .collect::<Vec<_>>();
265 let detached_count = targets.len();
266 for request_key in targets {
267 self.detach_request_from_provider(&request_key);
268 }
269 let retained_effect_count = self.effects.len();
270 self.effects
271 .retain(|effect| &effect.provider_id != provider_id);
272 let purged_effect_count = retained_effect_count - self.effects.len();
273 if detached_count > 0 || purged_effect_count > 0 {
274 self.bump_revision();
275 }
276 Ok(detached_count)
277 }
278
279 fn set_grant(&mut self, grant: PolicyGrant) -> Result<(), ToolEngineError> {
280 let Some(provider) = self.providers.get(&grant.key.provider_id) else {
281 return Err(ToolEngineError::UnknownPolicyProvider {
282 provider_id: grant.key.provider_id,
283 });
284 };
285 if !provider.has_capability(&grant.key.capability_id) {
286 return Err(ToolEngineError::UnknownPolicyCapability {
287 provider_id: grant.key.provider_id,
288 capability_id: grant.key.capability_id,
289 });
290 }
291 if let CapabilityOwner::Consumer(owner) = &provider.owner {
292 if owner != &grant.key.consumer_id {
293 return Err(ToolEngineError::ProviderOwnerMismatch {
294 provider_id: provider.id.clone(),
295 owner: owner.clone(),
296 requested_consumer: grant.key.consumer_id,
297 });
298 }
299 }
300 let Some(current_generation) = self.generations.get(&grant.key.instance_id).copied() else {
301 return Err(ToolEngineError::UnknownPolicyInstance {
302 instance_id: grant.key.instance_id,
303 });
304 };
305 if self.instance_states.get(&grant.key.instance_id) != Some(&ToolInstanceState::Active) {
306 return Err(ToolEngineError::InactivePolicyInstance {
307 instance_id: grant.key.instance_id,
308 });
309 }
310 if current_generation != grant.key.generation {
311 return Err(ToolEngineError::PolicyGenerationMismatch {
312 instance_id: grant.key.instance_id,
313 current: current_generation,
314 requested: grant.key.generation,
315 });
316 }
317 let previous = self.grants.get(&grant.key).copied();
318 if previous == Some(grant.mode) {
319 return Ok(());
320 }
321 if previous.is_none() && self.grants.len() >= TOOL_POLICIES_MAX {
322 return Err(ToolEngineError::PolicyCapacityExceeded);
323 }
324 if previous.is_some() {
325 let targets = self.active_requests_for_key(&grant.key);
326 for request_id in targets {
327 self.revoke_request(&request_id);
328 }
329 }
330 self.grants.insert(grant.key.clone(), grant.mode);
331 self.bump_revision();
332 self.emit_audit(None, ToolAuditEventKind::GrantSet { grant });
333 Ok(())
334 }
335
336 fn revoke_grant(&mut self, key: &PolicyKey) -> Result<bool, ToolEngineError> {
337 if !self.grants.contains_key(key) {
338 return Ok(false);
339 }
340 let targets = self.active_requests_for_key(key);
341 self.grants.remove(key);
342 for request_id in targets {
343 self.revoke_request(&request_id);
344 }
345 self.bump_revision();
346 self.emit_audit(None, ToolAuditEventKind::GrantRevoked { key: key.clone() });
347 Ok(true)
348 }
349
350 pub fn set_generation(
351 &mut self,
352 instance_id: AgentInstanceId,
353 generation: SessionGeneration,
354 ) -> Result<(), ToolEngineError> {
355 let previous = self.generations.get(&instance_id).copied();
356 if let Some(current) = previous {
357 if generation.0 < current.0 {
358 return Err(ToolEngineError::GenerationRegressed {
359 instance_id,
360 current,
361 requested: generation,
362 });
363 }
364 if generation == current {
365 return Ok(());
366 }
367 }
368
369 let stale_grants = self
370 .grants
371 .keys()
372 .filter(|key| key.instance_id == instance_id && key.generation != generation)
373 .cloned()
374 .collect::<Vec<_>>();
375 let purged_grant_count = stale_grants.len();
376 for key in stale_grants {
377 self.grants.remove(&key);
378 }
379 let targets = self
380 .requests
381 .iter()
382 .filter_map(|(request_id, state)| {
383 (state.request.instance_id == instance_id
384 && state.request.generation != generation
385 && !state.snapshot.status.is_terminal())
386 .then_some(request_id.clone())
387 })
388 .collect::<Vec<_>>();
389 self.generations.insert(instance_id, generation);
390 self.instance_states
391 .entry(instance_id)
392 .or_insert(ToolInstanceState::Active);
393 for request_id in targets {
394 self.supersede_request(&request_id, generation);
395 }
396 self.bump_revision();
397 self.emit_audit(
398 None,
399 ToolAuditEventKind::GenerationAdvanced {
400 instance_id,
401 previous,
402 current: generation,
403 purged_grant_count,
404 },
405 );
406 Ok(())
407 }
408
409 pub fn set_instance_state(
410 &mut self,
411 instance_id: AgentInstanceId,
412 generation: SessionGeneration,
413 state: ToolInstanceState,
414 ) -> Result<(), ToolEngineError> {
415 let Some(current_generation) = self.generations.get(&instance_id).copied() else {
416 return Err(ToolEngineError::UnknownPolicyInstance { instance_id });
417 };
418 if current_generation != generation {
419 return Err(ToolEngineError::PolicyGenerationMismatch {
420 instance_id,
421 current: current_generation,
422 requested: generation,
423 });
424 }
425 if self.instance_states.get(&instance_id) == Some(&state) {
426 return Ok(());
427 }
428 let mut purged_grant_count = 0;
429 if state == ToolInstanceState::Inactive {
430 purged_grant_count = self.purge_instance_grants(instance_id);
431 let targets = self
432 .requests
433 .iter()
434 .filter_map(|(request_key, request_state)| {
435 (request_state.request.instance_id == instance_id
436 && !request_state.snapshot.status.is_terminal())
437 .then_some(request_key.clone())
438 })
439 .collect::<Vec<_>>();
440 self.instance_states.insert(instance_id, state);
441 for request_key in targets {
442 self.close_request(&request_key, RequestCloseKind::Instance);
443 }
444 } else {
445 self.instance_states.insert(instance_id, state);
446 }
447 self.bump_revision();
448 self.emit_audit(
449 None,
450 ToolAuditEventKind::InstanceStateChanged {
451 instance_id,
452 generation,
453 state,
454 purged_grant_count,
455 },
456 );
457 Ok(())
458 }
459
460 pub fn remove_instance(
461 &mut self,
462 instance_id: AgentInstanceId,
463 ) -> Result<bool, ToolEngineError> {
464 let Some(previous_generation) = self.generations.get(&instance_id).copied() else {
465 return Ok(false);
466 };
467 let purged_grant_count = self.purge_instance_grants(instance_id);
468 let targets = self
469 .requests
470 .iter()
471 .filter_map(|(request_key, state)| {
472 (state.request.instance_id == instance_id && !state.snapshot.status.is_terminal())
473 .then_some(request_key.clone())
474 })
475 .collect::<Vec<_>>();
476 for request_key in targets {
477 self.close_request(&request_key, RequestCloseKind::Instance);
478 }
479 self.generations.remove(&instance_id);
480 self.instance_states.remove(&instance_id);
481 self.bump_revision();
482 self.emit_audit(
483 None,
484 ToolAuditEventKind::InstanceRemoved {
485 instance_id,
486 previous_generation,
487 purged_grant_count,
488 },
489 );
490 Ok(true)
491 }
492
493 fn close_client(
494 &mut self,
495 consumer_id: &ConsumerId,
496 actor_id: &ToolActorId,
497 ) -> ToolAuthorityOutcome {
498 let grant_keys = self
499 .grants
500 .keys()
501 .filter(|key| &key.consumer_id == consumer_id && &key.actor_id == actor_id)
502 .cloned()
503 .collect::<Vec<_>>();
504 let purged_grant_count = grant_keys.len();
505 for key in grant_keys {
506 self.grants.remove(&key);
507 }
508 let targets = self
509 .requests
510 .iter()
511 .filter_map(|(request_key, state)| {
512 (&request_key.consumer_id == consumer_id
513 && &request_key.actor_id == actor_id
514 && !state.snapshot.status.is_terminal())
515 .then_some(request_key.clone())
516 })
517 .collect::<Vec<_>>();
518 let closed_request_count = targets.len();
519 for request_key in targets {
520 self.close_request(&request_key, RequestCloseKind::Client);
521 }
522 self.bump_revision();
523 self.emit_audit(
524 None,
525 ToolAuditEventKind::ClientClosed {
526 consumer_id: consumer_id.clone(),
527 actor_id: actor_id.clone(),
528 purged_grant_count,
529 closed_request_count,
530 },
531 );
532 ToolAuthorityOutcome::ClientClosed {
533 purged_grant_count,
534 closed_request_count,
535 }
536 }
537
538 pub fn apply_authority(
541 &mut self,
542 envelope: ToolAuthorityEnvelope,
543 ) -> Result<ToolAuthorityOutcome, ToolEngineError> {
544 envelope.validate()?;
545 if envelope.sequence <= self.last_authority_sequence {
546 return Err(ToolEngineError::AuthoritySequenceRegressed {
547 current: self.last_authority_sequence,
548 requested: envelope.sequence,
549 });
550 }
551 let outcome = match envelope.command {
552 ToolAuthorityCommand::SetGrant { grant } => {
553 self.set_grant(grant)?;
554 ToolAuthorityOutcome::GrantSet
555 }
556 ToolAuthorityCommand::RevokeGrant { key } => ToolAuthorityOutcome::GrantRevoked {
557 existed: self.revoke_grant(&key)?,
558 },
559 ToolAuthorityCommand::ResolveApproval { resolution } => {
560 match self.resolve_approval(resolution.clone()) {
561 Ok(()) => ToolAuthorityOutcome::ApprovalResolved,
562 Err(ToolEngineError::ApprovalDeadlineElapsed { .. }) => {
563 ToolAuthorityOutcome::ApprovalExpired {
564 request_key: resolution.request_key,
565 accepted_sequence: resolution.accepted_sequence,
566 }
567 }
568 Err(error) => return Err(error),
569 }
570 }
571 ToolAuthorityCommand::CloseClient {
572 consumer_id,
573 actor_id,
574 } => self.close_client(&consumer_id, &actor_id),
575 };
576 self.last_authority_sequence = envelope.sequence;
577 Ok(outcome)
578 }
579
580 pub fn request(
581 &mut self,
582 envelope: ConsumerBoundCapabilityRequest,
583 ) -> Result<PolicyDecision, ToolEngineError> {
584 envelope.validate(self.current_tick)?;
585 let request_key = envelope.key();
586 let reuses_terminal_key = self
587 .requests
588 .get(&request_key)
589 .map(|state| state.snapshot.status.is_terminal())
590 .unwrap_or(false);
591 if self.requests.contains_key(&request_key) && !reuses_terminal_key {
592 return Err(ToolEngineError::DuplicateRequest { request_key });
593 }
594 let request = AcceptedRequest::from(envelope);
595 let decision = self.evaluate_policy(&request);
596 self.ensure_correlation_capacity(decision == PolicyDecision::Allow)?;
597 if decision == PolicyDecision::Allow {
598 self.ensure_effect_capacity(1)?;
599 }
600 if !matches!(decision, PolicyDecision::Deny(_)) {
601 self.ensure_client_request_capacity(&request.key)?;
602 }
603 if !reuses_terminal_key {
604 self.ensure_request_capacity()?;
605 }
606 let accepted_sequence = self.allocate_request_sequence();
607 let operation_id = (decision == PolicyDecision::Allow).then(|| self.allocate_operation());
608 let status = match decision {
609 PolicyDecision::Deny(reason) => CapabilityRequestStatus::Denied { reason },
610 PolicyDecision::RequireApproval => CapabilityRequestStatus::AwaitingApproval,
611 PolicyDecision::Allow => CapabilityRequestStatus::Dispatched {
612 operation_id: operation_id.expect("operation allocated for allowed request"),
613 },
614 };
615 let snapshot = CapabilityRequestSnapshot {
616 key: request.key.clone(),
617 accepted_sequence,
618 accepted_at_tick: self.current_tick,
619 instance_id: request.instance_id,
620 generation: request.generation,
621 provider_id: request.provider_id.clone(),
622 capability_id: request.capability_id.clone(),
623 resource_scope_id: request.resource_scope_id.clone(),
624 approval_summary: request.approval_summary.clone(),
625 approval_summary_bytes: request.approval_summary.len(),
626 deadline_tick: request.deadline_tick,
627 payload_bytes: request.payload.len(),
628 policy_decision: decision,
629 status,
630 };
631 let subject = subject_for(&request, accepted_sequence);
632 let payload_bytes = request.payload.len();
633 let mut stored_request = request.clone();
634 if decision != PolicyDecision::RequireApproval {
635 stored_request.payload.clear();
636 }
637 if reuses_terminal_key {
638 self.requests.remove(&request.key);
639 }
640 self.requests.insert(
641 request.key.clone(),
642 RequestState {
643 request: stored_request,
644 snapshot,
645 },
646 );
647 self.bump_revision();
648 self.emit_audit(
649 Some(subject.clone()),
650 ToolAuditEventKind::RequestEvaluated {
651 decision,
652 payload_bytes,
653 },
654 );
655 match (decision, operation_id) {
656 (PolicyDecision::Deny(reason), None) => self.push_terminal_completion(
657 &request,
658 accepted_sequence,
659 None,
660 CapabilityTerminalOutcome::PolicyDenied { reason },
661 ),
662 (PolicyDecision::Allow, Some(operation_id)) => {
663 self.emit_invoke(&request, operation_id);
664 self.emit_audit(
665 Some(subject),
666 ToolAuditEventKind::InvocationDispatched { operation_id },
667 );
668 }
669 _ => {}
670 }
671 Ok(decision)
672 }
673
674 fn resolve_approval(&mut self, resolution: ApprovalResolution) -> Result<(), ToolEngineError> {
675 let Some(state) = self.requests.get(&resolution.request_key) else {
676 return Err(ToolEngineError::UnknownRequest {
677 request_key: resolution.request_key,
678 });
679 };
680 if resolution.accepted_sequence != state.snapshot.accepted_sequence {
681 return Err(ToolEngineError::ApprovalNonceMismatch {
682 request_key: resolution.request_key,
683 expected: state.snapshot.accepted_sequence,
684 actual: resolution.accepted_sequence,
685 });
686 }
687 let Some(current) = self.generations.get(&state.request.instance_id).copied() else {
688 return Err(ToolEngineError::UnknownPolicyInstance {
689 instance_id: state.request.instance_id,
690 });
691 };
692 if self.instance_states.get(&state.request.instance_id) != Some(&ToolInstanceState::Active)
693 {
694 return Err(ToolEngineError::InactivePolicyInstance {
695 instance_id: state.request.instance_id,
696 });
697 }
698 if resolution.generation != current {
699 return Err(ToolEngineError::ApprovalGenerationStale {
700 current,
701 actual: resolution.generation,
702 });
703 }
704 if resolution.instance_id != state.request.instance_id
705 || resolution.generation != state.request.generation
706 {
707 return Err(ToolEngineError::ApprovalScopeMismatch {
708 request_key: resolution.request_key,
709 });
710 }
711 if !matches!(
712 state.snapshot.status,
713 CapabilityRequestStatus::AwaitingApproval
714 ) {
715 return Err(ToolEngineError::RequestNotAwaitingApproval {
716 request_key: resolution.request_key,
717 });
718 }
719 if state.request.deadline_tick <= self.current_tick {
720 self.timeout_request(&resolution.request_key);
721 self.bump_revision();
722 return Err(ToolEngineError::ApprovalDeadlineElapsed {
723 request_key: resolution.request_key,
724 });
725 }
726 if resolution.decision == ApprovalDecision::ApproveOnce {
727 self.ensure_operation_capacity()?;
728 self.ensure_effect_capacity(1)?;
729 }
730
731 let operation_id = (resolution.decision == ApprovalDecision::ApproveOnce)
732 .then(|| self.allocate_operation());
733 let (request, subject) = {
734 let state = self
735 .requests
736 .get_mut(&resolution.request_key)
737 .expect("approval request checked above");
738 let request = state.request.clone();
739 state.snapshot.status = match operation_id {
740 Some(operation_id) => CapabilityRequestStatus::Dispatched { operation_id },
741 None => CapabilityRequestStatus::ApprovalDenied,
742 };
743 state.request.payload.clear();
744 (
745 request,
746 subject_for(&state.request, state.snapshot.accepted_sequence),
747 )
748 };
749 self.bump_revision();
750 self.emit_audit(
751 Some(subject.clone()),
752 ToolAuditEventKind::ApprovalResolved {
753 decision: resolution.decision,
754 },
755 );
756 if let Some(operation_id) = operation_id {
757 self.emit_invoke(&request, operation_id);
758 self.emit_audit(
759 Some(subject),
760 ToolAuditEventKind::InvocationDispatched { operation_id },
761 );
762 } else {
763 self.push_terminal_completion(
764 &request,
765 resolution.accepted_sequence,
766 None,
767 CapabilityTerminalOutcome::ApprovalDenied,
768 );
769 }
770 Ok(())
771 }
772
773 pub fn advance_time(&mut self, tick: u64) -> Result<(), ToolEngineError> {
774 if tick < self.current_tick {
775 return Err(ToolEngineError::ClockRegressed {
776 current_tick: self.current_tick,
777 requested_tick: tick,
778 });
779 }
780 if tick == self.current_tick {
781 return Ok(());
782 }
783 let expired = self
784 .requests
785 .iter()
786 .filter_map(|(request_id, state)| {
787 (!state.snapshot.status.is_terminal() && state.request.deadline_tick <= tick)
788 .then_some(request_id.clone())
789 })
790 .collect::<Vec<_>>();
791 self.current_tick = tick;
792 for request_id in expired {
793 self.timeout_request(&request_id);
794 }
795 self.bump_revision();
796 Ok(())
797 }
798
799 pub fn apply_observation(
802 &mut self,
803 envelope: CapabilityObservationEnvelope,
804 ) -> Result<CapabilityObservationDisposition, ToolEngineError> {
805 if envelope.operation_id.0 == 0 {
806 return Err(ToolValidationError::ZeroIdentifier {
807 field: "tool operation id",
808 }
809 .into());
810 }
811 let Some(state) = self.requests.get(&envelope.request_key) else {
812 return Ok(self.ignore_observation(
813 &envelope,
814 None,
815 ObservationIgnoredReason::UnknownRequest,
816 ));
817 };
818 let subject = subject_for(&state.request, state.snapshot.accepted_sequence);
819 let current_generation = self.generations.get(&state.request.instance_id).copied();
820 if current_generation != Some(envelope.generation)
821 || envelope.generation != state.request.generation
822 {
823 return Ok(self.ignore_observation(
824 &envelope,
825 Some(subject),
826 ObservationIgnoredReason::StaleGeneration,
827 ));
828 }
829 if envelope.instance_id != state.request.instance_id {
830 return Ok(self.ignore_observation(
831 &envelope,
832 Some(subject),
833 ObservationIgnoredReason::InstanceMismatch,
834 ));
835 }
836 if envelope.provider_id != state.request.provider_id {
837 return Ok(self.ignore_observation(
838 &envelope,
839 Some(subject),
840 ObservationIgnoredReason::ProviderMismatch,
841 ));
842 }
843 let operation_id = match state.snapshot.status {
844 CapabilityRequestStatus::Dispatched { operation_id } => operation_id,
845 _ => {
846 return Ok(self.ignore_observation(
847 &envelope,
848 Some(subject),
849 ObservationIgnoredReason::RequestNotDispatched,
850 ));
851 }
852 };
853 if envelope.operation_id != operation_id {
854 return Ok(self.ignore_observation(
855 &envelope,
856 Some(subject),
857 ObservationIgnoredReason::OperationMismatch,
858 ));
859 }
860 if state.request.deadline_tick <= self.current_tick {
861 self.timeout_request(&envelope.request_key);
862 self.bump_revision();
863 self.emit_audit(
864 Some(subject),
865 ToolAuditEventKind::ObservationIgnored {
866 operation_id,
867 reason: ObservationIgnoredReason::DeadlineElapsed,
868 },
869 );
870 return Ok(CapabilityObservationDisposition::Ignored {
871 reason: ObservationIgnoredReason::DeadlineElapsed,
872 });
873 }
874
875 if envelope.observation.validate().is_err() {
876 let failure = ToolFailure::provider_contract_violation();
877 let state = self
878 .requests
879 .get_mut(&envelope.request_key)
880 .expect("observation request checked above");
881 state.snapshot.status = CapabilityRequestStatus::Failed {
882 operation_id,
883 failure: failure.clone(),
884 };
885 state.request.payload.clear();
886 let accepted_sequence = state.snapshot.accepted_sequence;
887 let request = state.request.clone();
888 self.bump_revision();
889 self.emit_audit(
890 Some(subject),
891 ToolAuditEventKind::InvocationFailed {
892 operation_id,
893 failure_kind: failure.kind,
894 },
895 );
896 self.push_terminal_completion(
897 &request,
898 accepted_sequence,
899 Some(operation_id),
900 CapabilityTerminalOutcome::Failed { failure },
901 );
902 return Ok(CapabilityObservationDisposition::Applied);
903 }
904
905 let event = match envelope.observation {
906 CapabilityObservation::Succeeded { result } => {
907 let event = ToolAuditEventKind::InvocationSucceeded {
908 operation_id,
909 result_bytes: result.metadata.byte_len,
910 truncated: result.metadata.truncated,
911 };
912 let metadata = result.metadata.clone();
913 let state = self
914 .requests
915 .get_mut(&envelope.request_key)
916 .expect("observation request checked above");
917 state.snapshot.status = CapabilityRequestStatus::Succeeded {
918 operation_id,
919 result: metadata,
920 };
921 state.request.payload.clear();
922 let accepted_sequence = state.snapshot.accepted_sequence;
923 let request = state.request.clone();
924 self.push_terminal_completion(
925 &request,
926 accepted_sequence,
927 Some(operation_id),
928 CapabilityTerminalOutcome::Succeeded { result },
929 );
930 event
931 }
932 CapabilityObservation::Failed { failure } => {
933 let event = ToolAuditEventKind::InvocationFailed {
934 operation_id,
935 failure_kind: failure.kind,
936 };
937 let state = self
938 .requests
939 .get_mut(&envelope.request_key)
940 .expect("observation request checked above");
941 state.snapshot.status = CapabilityRequestStatus::Failed {
942 operation_id,
943 failure: failure.clone(),
944 };
945 state.request.payload.clear();
946 let accepted_sequence = state.snapshot.accepted_sequence;
947 let request = state.request.clone();
948 self.push_terminal_completion(
949 &request,
950 accepted_sequence,
951 Some(operation_id),
952 CapabilityTerminalOutcome::Failed { failure },
953 );
954 event
955 }
956 };
957 self.bump_revision();
958 self.emit_audit(Some(subject), event);
959 Ok(CapabilityObservationDisposition::Applied)
960 }
961
962 pub fn snapshot(&self) -> ToolEngineSnapshot {
963 ToolEngineSnapshot {
964 revision: self.revision,
965 current_tick: self.current_tick,
966 generations: self
967 .generations
968 .iter()
969 .map(|(instance_id, generation)| (*instance_id, *generation))
970 .collect(),
971 instance_states: self
972 .instance_states
973 .iter()
974 .map(|(instance_id, state)| (*instance_id, *state))
975 .collect(),
976 providers: self.providers.values().cloned().collect(),
977 grants: self
978 .grants
979 .iter()
980 .map(|(key, mode)| PolicyGrant {
981 key: key.clone(),
982 mode: *mode,
983 })
984 .collect(),
985 requests: self
986 .requests
987 .values()
988 .map(|state| state.snapshot.clone())
989 .collect(),
990 audit_events: self.audit_events.iter().cloned().collect(),
991 dropped_audit_events: self.dropped_audit_events,
992 revision_overflow_count: self.revision_overflow_count,
993 next_completion_sequence: self.next_completion_sequence,
994 dropped_completions: self.dropped_completions,
995 effect_sequence_exhausted: self.effect_sequence_exhausted,
996 completion_sequence_exhausted: self.completion_sequence_exhausted,
997 audit_sequence_exhausted: self.audit_sequence_exhausted,
998 }
999 }
1000
1001 pub fn request_snapshot(
1002 &self,
1003 request_key: &CapabilityRequestKey,
1004 ) -> Option<&CapabilityRequestSnapshot> {
1005 self.requests.get(request_key).map(|state| &state.snapshot)
1006 }
1007
1008 pub fn instance_ids(&self) -> impl Iterator<Item = AgentInstanceId> + '_ {
1009 self.generations.keys().copied()
1010 }
1011
1012 pub fn drain_effects(&mut self) -> Vec<CapabilityEffectEnvelope> {
1013 std::mem::take(&mut self.effects)
1014 }
1015
1016 pub fn drain_completions(&mut self) -> CapabilityCompletionBatch {
1020 let dropped_since_last_drain = std::mem::take(&mut self.dropped_completions_since_drain);
1021 CapabilityCompletionBatch {
1022 completions: std::mem::take(&mut self.completions),
1023 dropped_since_last_drain,
1024 total_dropped: self.dropped_completions,
1025 next_sequence: self.next_completion_sequence,
1026 sequence_exhausted: self.completion_sequence_exhausted,
1027 }
1028 }
1029
1030 fn evaluate_policy(&self, request: &AcceptedRequest) -> PolicyDecision {
1031 let Some(current_generation) = self.generations.get(&request.instance_id).copied() else {
1032 return PolicyDecision::Deny(PolicyDenial::UnknownInstance);
1033 };
1034 if self.instance_states.get(&request.instance_id) != Some(&ToolInstanceState::Active) {
1035 return PolicyDecision::Deny(PolicyDenial::InactiveInstance);
1036 }
1037 if current_generation != request.generation {
1038 return PolicyDecision::Deny(PolicyDenial::StaleGeneration {
1039 current: current_generation,
1040 });
1041 }
1042 let Some(provider) = self.providers.get(&request.provider_id) else {
1043 return PolicyDecision::Deny(PolicyDenial::UnknownProvider);
1044 };
1045 if matches!(
1046 &provider.owner,
1047 CapabilityOwner::Consumer(owner) if owner != &request.key.consumer_id
1048 ) {
1049 return PolicyDecision::Deny(PolicyDenial::ProviderOwnerMismatch);
1050 }
1051 if !provider.has_capability(&request.capability_id) {
1052 return PolicyDecision::Deny(PolicyDenial::UnknownCapability);
1053 }
1054 let key = PolicyKey {
1055 consumer_id: request.key.consumer_id.clone(),
1056 actor_id: request.key.actor_id.clone(),
1057 instance_id: request.instance_id,
1058 generation: request.generation,
1059 provider_id: request.provider_id.clone(),
1060 capability_id: request.capability_id.clone(),
1061 resource_scope_id: request.resource_scope_id.clone(),
1062 };
1063 match self.grants.get(&key) {
1064 Some(GrantMode::Allow) => PolicyDecision::Allow,
1065 Some(GrantMode::RequireApproval) => PolicyDecision::RequireApproval,
1066 None => PolicyDecision::Deny(PolicyDenial::MissingGrant),
1067 }
1068 }
1069
1070 fn active_requests_for_key(&self, key: &PolicyKey) -> Vec<CapabilityRequestKey> {
1071 self.requests
1072 .iter()
1073 .filter_map(|(request_id, state)| {
1074 (request_matches_key(&state.request, key) && !state.snapshot.status.is_terminal())
1075 .then_some(request_id.clone())
1076 })
1077 .collect()
1078 }
1079
1080 fn detach_request_from_provider(&mut self, request_key: &CapabilityRequestKey) {
1081 let (request, subject, accepted_sequence, operation_id) = {
1082 let state = self
1083 .requests
1084 .get_mut(request_key)
1085 .expect("provider detach target collected from request map");
1086 let request = state.request.clone();
1087 let operation_id = match state.snapshot.status {
1088 CapabilityRequestStatus::Dispatched { operation_id } => Some(operation_id),
1089 _ => None,
1090 };
1091 state.request.payload.clear();
1092 (
1093 request,
1094 subject_for(&state.request, state.snapshot.accepted_sequence),
1095 state.snapshot.accepted_sequence,
1096 operation_id,
1097 )
1098 };
1099 let cancellation = match operation_id {
1100 Some(operation_id) => {
1101 let queued_invoke = self.effects.iter().position(|effect| {
1102 effect.request_key == request.key
1103 && effect.operation_id == operation_id
1104 && matches!(&effect.effect, CapabilityEffect::Invoke { .. })
1105 });
1106 if let Some(position) = queued_invoke {
1107 self.effects.remove(position);
1108 CancellationDisposition::QueuedInvokeRemoved
1109 } else {
1110 CancellationDisposition::ProviderDetachedUnconfirmed
1111 }
1112 }
1113 None => CancellationDisposition::NotRequired,
1114 };
1115 self.requests
1116 .get_mut(request_key)
1117 .expect("provider detach target remains in request map")
1118 .snapshot
1119 .status = CapabilityRequestStatus::ProviderDetached {
1120 operation_id,
1121 cancellation,
1122 };
1123 self.emit_audit(
1124 Some(subject),
1125 ToolAuditEventKind::RequestProviderDetached {
1126 operation_id,
1127 cancellation,
1128 },
1129 );
1130 self.push_terminal_completion(
1131 &request,
1132 accepted_sequence,
1133 operation_id,
1134 CapabilityTerminalOutcome::ProviderDetached { cancellation },
1135 );
1136 }
1137
1138 fn supersede_request(
1139 &mut self,
1140 request_key: &CapabilityRequestKey,
1141 current_generation: SessionGeneration,
1142 ) {
1143 let (request, subject, accepted_sequence, operation_id) = {
1144 let state = self
1145 .requests
1146 .get_mut(request_key)
1147 .expect("supersede target collected from request map");
1148 let request = state.request.clone();
1149 let operation_id = match state.snapshot.status {
1150 CapabilityRequestStatus::Dispatched { operation_id } => Some(operation_id),
1151 _ => None,
1152 };
1153 state.request.payload.clear();
1154 (
1155 request,
1156 subject_for(&state.request, state.snapshot.accepted_sequence),
1157 state.snapshot.accepted_sequence,
1158 operation_id,
1159 )
1160 };
1161 let cancellation = match operation_id {
1162 Some(operation_id) => self.cancel_best_effort(
1163 &request,
1164 operation_id,
1165 InvocationCancelReason::GenerationSuperseded,
1166 ),
1167 None => CancellationDisposition::NotRequired,
1168 };
1169 self.requests
1170 .get_mut(request_key)
1171 .expect("supersede target remains in request map")
1172 .snapshot
1173 .status = CapabilityRequestStatus::Superseded {
1174 operation_id,
1175 cancellation,
1176 current_generation,
1177 };
1178 self.emit_audit(
1179 Some(subject),
1180 ToolAuditEventKind::RequestSuperseded {
1181 operation_id,
1182 cancellation,
1183 current_generation,
1184 },
1185 );
1186 self.push_terminal_completion(
1187 &request,
1188 accepted_sequence,
1189 operation_id,
1190 CapabilityTerminalOutcome::Superseded {
1191 current_generation,
1192 cancellation,
1193 },
1194 );
1195 }
1196
1197 fn timeout_request(&mut self, request_key: &CapabilityRequestKey) {
1198 let (request, subject, accepted_sequence, operation_id) = {
1199 let state = self
1200 .requests
1201 .get_mut(request_key)
1202 .expect("timeout target collected from request map");
1203 let request = state.request.clone();
1204 let operation_id = match state.snapshot.status {
1205 CapabilityRequestStatus::Dispatched { operation_id } => Some(operation_id),
1206 _ => None,
1207 };
1208 state.request.payload.clear();
1209 (
1210 request,
1211 subject_for(&state.request, state.snapshot.accepted_sequence),
1212 state.snapshot.accepted_sequence,
1213 operation_id,
1214 )
1215 };
1216 let cancellation = match operation_id {
1217 Some(operation_id) => self.cancel_best_effort(
1218 &request,
1219 operation_id,
1220 InvocationCancelReason::DeadlineElapsed,
1221 ),
1222 None => CancellationDisposition::NotRequired,
1223 };
1224 self.requests
1225 .get_mut(request_key)
1226 .expect("timeout target remains in request map")
1227 .snapshot
1228 .status = CapabilityRequestStatus::TimedOut {
1229 operation_id,
1230 cancellation,
1231 };
1232 self.emit_audit(
1233 Some(subject),
1234 ToolAuditEventKind::RequestTimedOut {
1235 operation_id,
1236 cancellation,
1237 },
1238 );
1239 self.push_terminal_completion(
1240 &request,
1241 accepted_sequence,
1242 operation_id,
1243 CapabilityTerminalOutcome::TimedOut { cancellation },
1244 );
1245 }
1246
1247 fn revoke_request(&mut self, request_key: &CapabilityRequestKey) {
1248 let (request, subject, operation_id) = {
1249 let state = self
1250 .requests
1251 .get_mut(request_key)
1252 .expect("revocation target collected from request map");
1253 let request = state.request.clone();
1254 let operation_id = match state.snapshot.status {
1255 CapabilityRequestStatus::Dispatched { operation_id } => Some(operation_id),
1256 _ => None,
1257 };
1258 state.request.payload.clear();
1259 (
1260 request,
1261 subject_for(&state.request, state.snapshot.accepted_sequence),
1262 operation_id,
1263 )
1264 };
1265 let cancellation = match operation_id {
1266 Some(operation_id) => self.cancel_best_effort(
1267 &request,
1268 operation_id,
1269 InvocationCancelReason::GrantRevoked,
1270 ),
1271 None => CancellationDisposition::NotRequired,
1272 };
1273 self.requests
1274 .get_mut(request_key)
1275 .expect("revocation target remains in request map")
1276 .snapshot
1277 .status = CapabilityRequestStatus::GrantRevoked {
1278 operation_id,
1279 cancellation,
1280 };
1281 self.emit_audit(
1282 Some(subject),
1283 ToolAuditEventKind::RequestGrantRevoked {
1284 operation_id,
1285 cancellation,
1286 },
1287 );
1288 self.push_terminal_completion(
1289 &request,
1290 self.requests[request_key].snapshot.accepted_sequence,
1291 operation_id,
1292 CapabilityTerminalOutcome::GrantRevoked { cancellation },
1293 );
1294 }
1295
1296 fn close_request(&mut self, request_key: &CapabilityRequestKey, kind: RequestCloseKind) {
1297 let (request, subject, accepted_sequence, operation_id) = {
1298 let state = self
1299 .requests
1300 .get_mut(request_key)
1301 .expect("close target collected from request map");
1302 let request = state.request.clone();
1303 let operation_id = match state.snapshot.status {
1304 CapabilityRequestStatus::Dispatched { operation_id } => Some(operation_id),
1305 _ => None,
1306 };
1307 state.request.payload.clear();
1308 (
1309 request,
1310 subject_for(&state.request, state.snapshot.accepted_sequence),
1311 state.snapshot.accepted_sequence,
1312 operation_id,
1313 )
1314 };
1315 let reason = match kind {
1316 RequestCloseKind::Instance => InvocationCancelReason::InstanceClosed,
1317 RequestCloseKind::Client => InvocationCancelReason::ClientClosed,
1318 };
1319 let cancellation = match operation_id {
1320 Some(operation_id) => self.cancel_best_effort(&request, operation_id, reason),
1321 None => CancellationDisposition::NotRequired,
1322 };
1323 let (status, outcome, audit) = match kind {
1324 RequestCloseKind::Instance => (
1325 CapabilityRequestStatus::InstanceClosed {
1326 operation_id,
1327 cancellation,
1328 },
1329 CapabilityTerminalOutcome::InstanceClosed { cancellation },
1330 ToolAuditEventKind::RequestInstanceClosed {
1331 operation_id,
1332 cancellation,
1333 },
1334 ),
1335 RequestCloseKind::Client => (
1336 CapabilityRequestStatus::ClientClosed {
1337 operation_id,
1338 cancellation,
1339 },
1340 CapabilityTerminalOutcome::ClientClosed { cancellation },
1341 ToolAuditEventKind::RequestClientClosed {
1342 operation_id,
1343 cancellation,
1344 },
1345 ),
1346 };
1347 self.requests
1348 .get_mut(request_key)
1349 .expect("close target remains in request map")
1350 .snapshot
1351 .status = status;
1352 self.emit_audit(Some(subject), audit);
1353 self.push_terminal_completion(&request, accepted_sequence, operation_id, outcome);
1354 }
1355
1356 fn purge_instance_grants(&mut self, instance_id: AgentInstanceId) -> usize {
1357 let keys = self
1358 .grants
1359 .keys()
1360 .filter(|key| key.instance_id == instance_id)
1361 .cloned()
1362 .collect::<Vec<_>>();
1363 let count = keys.len();
1364 for key in keys {
1365 self.grants.remove(&key);
1366 }
1367 count
1368 }
1369
1370 fn ignore_observation(
1371 &mut self,
1372 envelope: &CapabilityObservationEnvelope,
1373 subject: Option<ToolAuditSubject>,
1374 reason: ObservationIgnoredReason,
1375 ) -> CapabilityObservationDisposition {
1376 self.bump_revision();
1377 self.emit_audit(
1378 subject,
1379 ToolAuditEventKind::ObservationIgnored {
1380 operation_id: envelope.operation_id,
1381 reason,
1382 },
1383 );
1384 CapabilityObservationDisposition::Ignored { reason }
1385 }
1386
1387 fn emit_invoke(&mut self, request: &AcceptedRequest, operation_id: ToolOperationId) {
1388 let emitted = self.push_effect(
1389 request,
1390 operation_id,
1391 CapabilityEffect::Invoke {
1392 consumer_id: request.key.consumer_id.clone(),
1393 actor_id: request.key.actor_id.clone(),
1394 capability_id: request.capability_id.clone(),
1395 resource_scope_id: request.resource_scope_id.clone(),
1396 payload: request.payload.clone(),
1397 },
1398 );
1399 debug_assert!(emitted, "invoke effect sequence was preflighted");
1400 }
1401
1402 fn cancel_best_effort(
1403 &mut self,
1404 request: &AcceptedRequest,
1405 operation_id: ToolOperationId,
1406 reason: InvocationCancelReason,
1407 ) -> CancellationDisposition {
1408 let queued_invoke = self.effects.iter().position(|effect| {
1409 effect.request_key == request.key
1410 && effect.operation_id == operation_id
1411 && matches!(&effect.effect, CapabilityEffect::Invoke { .. })
1412 });
1413 if let Some(position) = queued_invoke {
1414 self.effects.remove(position);
1415 return CancellationDisposition::QueuedInvokeRemoved;
1416 }
1417 if self.effects.len() >= TOOL_EFFECTS_MAX {
1418 return CancellationDisposition::DroppedQueueFull;
1419 }
1420 if self.emit_cancel(request, operation_id, reason) {
1421 CancellationDisposition::CancelQueuedUnconfirmed
1422 } else {
1423 CancellationDisposition::DroppedSequenceExhausted
1424 }
1425 }
1426
1427 fn emit_cancel(
1428 &mut self,
1429 request: &AcceptedRequest,
1430 operation_id: ToolOperationId,
1431 reason: InvocationCancelReason,
1432 ) -> bool {
1433 self.push_effect(request, operation_id, CapabilityEffect::Cancel { reason })
1434 }
1435
1436 fn push_effect(
1437 &mut self,
1438 request: &AcceptedRequest,
1439 operation_id: ToolOperationId,
1440 effect: CapabilityEffect,
1441 ) -> bool {
1442 if self.effect_sequence_exhausted {
1443 return false;
1444 }
1445 let sequence = self.next_effect_sequence;
1446 if sequence == u64::MAX {
1447 self.effect_sequence_exhausted = true;
1448 } else {
1449 self.next_effect_sequence += 1;
1450 }
1451 self.effects.push(CapabilityEffectEnvelope {
1452 sequence,
1453 operation_id,
1454 request_key: request.key.clone(),
1455 instance_id: request.instance_id,
1456 generation: request.generation,
1457 provider_id: request.provider_id.clone(),
1458 deadline_tick: request.deadline_tick,
1459 effect,
1460 });
1461 true
1462 }
1463
1464 fn push_terminal_completion(
1465 &mut self,
1466 request: &AcceptedRequest,
1467 accepted_sequence: u64,
1468 operation_id: Option<ToolOperationId>,
1469 outcome: CapabilityTerminalOutcome,
1470 ) {
1471 if self.completion_sequence_exhausted {
1472 self.record_completion_drop(
1473 request,
1474 accepted_sequence,
1475 None,
1476 CompletionDropReason::SequenceExhausted,
1477 );
1478 return;
1479 }
1480 let sequence = self.next_completion_sequence;
1481 if sequence == u64::MAX {
1482 self.completion_sequence_exhausted = true;
1483 } else {
1484 self.next_completion_sequence += 1;
1485 }
1486 let completion = CapabilityCompletionEnvelope {
1487 sequence,
1488 accepted_sequence,
1489 operation_id,
1490 request_key: request.key.clone(),
1491 instance_id: request.instance_id,
1492 generation: request.generation,
1493 provider_id: request.provider_id.clone(),
1494 outcome,
1495 };
1496 if self.completions.len() >= TOOL_COMPLETIONS_MAX {
1497 self.record_completion_drop(
1498 request,
1499 accepted_sequence,
1500 Some(sequence),
1501 CompletionDropReason::QueueFull,
1502 );
1503 return;
1504 }
1505 self.completions.push(completion);
1506 }
1507
1508 fn record_completion_drop(
1509 &mut self,
1510 request: &AcceptedRequest,
1511 accepted_sequence: u64,
1512 completion_sequence: Option<u64>,
1513 reason: CompletionDropReason,
1514 ) {
1515 self.dropped_completions = self.dropped_completions.saturating_add(1);
1516 self.dropped_completions_since_drain =
1517 self.dropped_completions_since_drain.saturating_add(1);
1518 self.emit_audit(
1519 Some(subject_for(request, accepted_sequence)),
1520 ToolAuditEventKind::CompletionDropped {
1521 completion_sequence,
1522 reason,
1523 },
1524 );
1525 }
1526
1527 fn ensure_request_capacity(&mut self) -> Result<(), ToolEngineError> {
1528 if self.requests.len() < TOOL_REQUESTS_MAX {
1529 return Ok(());
1530 }
1531 let oldest_terminal = self
1532 .requests
1533 .iter()
1534 .filter(|(_, state)| state.snapshot.status.is_terminal())
1535 .min_by_key(|(_, state)| state.snapshot.accepted_sequence)
1536 .map(|(request_key, _)| request_key.clone());
1537 if let Some(request_key) = oldest_terminal {
1538 self.requests.remove(&request_key);
1539 return Ok(());
1540 }
1541 Err(ToolEngineError::RequestCapacityExceeded)
1542 }
1543
1544 fn ensure_effect_capacity(&self, additional: usize) -> Result<(), ToolEngineError> {
1545 if self.effect_sequence_exhausted {
1546 return Err(ToolEngineError::EffectSequenceExhausted);
1547 }
1548 if additional <= TOOL_EFFECTS_MAX.saturating_sub(self.effects.len()) {
1549 Ok(())
1550 } else {
1551 Err(ToolEngineError::EffectCapacityExceeded)
1552 }
1553 }
1554
1555 fn ensure_client_request_capacity(
1556 &self,
1557 request_key: &CapabilityRequestKey,
1558 ) -> Result<(), ToolEngineError> {
1559 let active_count = self
1560 .requests
1561 .iter()
1562 .filter(|(key, state)| {
1563 key.consumer_id == request_key.consumer_id
1564 && key.actor_id == request_key.actor_id
1565 && !state.snapshot.status.is_terminal()
1566 })
1567 .count();
1568 if active_count < TOOL_ACTIVE_REQUESTS_PER_CLIENT_MAX {
1569 Ok(())
1570 } else {
1571 Err(ToolEngineError::ClientRequestCapacityExceeded {
1572 consumer_id: request_key.consumer_id.clone(),
1573 actor_id: request_key.actor_id.clone(),
1574 max: TOOL_ACTIVE_REQUESTS_PER_CLIENT_MAX,
1575 })
1576 }
1577 }
1578
1579 fn ensure_correlation_capacity(&self, needs_operation: bool) -> Result<(), ToolEngineError> {
1580 if self.next_request_sequence == u64::MAX {
1581 return Err(ToolEngineError::CounterExhausted {
1582 counter: "tool request sequence",
1583 });
1584 }
1585 if needs_operation {
1586 self.ensure_operation_capacity()?;
1587 }
1588 Ok(())
1589 }
1590
1591 fn ensure_operation_capacity(&self) -> Result<(), ToolEngineError> {
1592 if self.next_operation_id == u64::MAX {
1593 Err(ToolEngineError::CounterExhausted {
1594 counter: "tool operation id",
1595 })
1596 } else {
1597 Ok(())
1598 }
1599 }
1600
1601 fn allocate_request_sequence(&mut self) -> u64 {
1602 let sequence = self.next_request_sequence;
1603 self.next_request_sequence += 1;
1604 sequence
1605 }
1606
1607 fn allocate_operation(&mut self) -> ToolOperationId {
1608 let operation_id = ToolOperationId(self.next_operation_id);
1609 self.next_operation_id += 1;
1610 operation_id
1611 }
1612
1613 fn bump_revision(&mut self) {
1614 if self.revision == u64::MAX {
1615 self.revision_overflow_count = self.revision_overflow_count.saturating_add(1);
1616 } else {
1617 self.revision += 1;
1618 }
1619 }
1620
1621 fn emit_audit(&mut self, subject: Option<ToolAuditSubject>, event: ToolAuditEventKind) {
1622 if self.audit_sequence_exhausted {
1623 self.dropped_audit_events = self.dropped_audit_events.saturating_add(1);
1624 return;
1625 }
1626 if self.audit_events.len() == TOOL_AUDIT_EVENTS_MAX {
1627 self.audit_events.pop_front();
1628 self.dropped_audit_events = self.dropped_audit_events.saturating_add(1);
1629 }
1630 let sequence = self.next_audit_sequence;
1631 if sequence == u64::MAX {
1632 self.audit_sequence_exhausted = true;
1633 } else {
1634 self.next_audit_sequence += 1;
1635 }
1636 self.audit_events.push_back(ToolAuditEvent {
1637 sequence,
1638 tick: self.current_tick,
1639 subject,
1640 event,
1641 });
1642 }
1643}
1644
1645impl Default for ToolEngine {
1646 fn default() -> Self {
1647 Self::new()
1648 }
1649}
1650
1651impl fmt::Debug for ToolEngine {
1652 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
1653 formatter
1654 .debug_struct("ToolEngine")
1655 .field("snapshot", &self.snapshot())
1656 .field("pending_effect_count", &self.effects.len())
1657 .field("pending_completion_count", &self.completions.len())
1658 .finish()
1659 }
1660}
1661
1662fn subject_for(request: &AcceptedRequest, accepted_sequence: u64) -> ToolAuditSubject {
1663 ToolAuditSubject {
1664 request_key: request.key.clone(),
1665 accepted_sequence,
1666 instance_id: request.instance_id,
1667 generation: request.generation,
1668 provider_id: request.provider_id.clone(),
1669 capability_id: request.capability_id.clone(),
1670 resource_scope_id: request.resource_scope_id.clone(),
1671 }
1672}
1673
1674fn request_matches_key(request: &AcceptedRequest, key: &PolicyKey) -> bool {
1675 request.key.consumer_id == key.consumer_id
1676 && request.key.actor_id == key.actor_id
1677 && request.instance_id == key.instance_id
1678 && request.generation == key.generation
1679 && request.provider_id == key.provider_id
1680 && request.capability_id == key.capability_id
1681 && request.resource_scope_id == key.resource_scope_id
1682}
1683
1684#[cfg(test)]
1685mod tests {
1686 use super::*;
1687
1688 fn instance() -> AgentInstanceId {
1689 AgentInstanceId(41)
1690 }
1691
1692 fn generation(value: u64) -> SessionGeneration {
1693 SessionGeneration(value)
1694 }
1695
1696 fn actor() -> ToolActorId {
1697 ToolActorId::new("consumer.agent").unwrap()
1698 }
1699
1700 fn consumer() -> ConsumerId {
1701 ConsumerId::new("station.test").unwrap()
1702 }
1703
1704 fn resource_scope() -> ResourceScopeId {
1705 ResourceScopeId::new("workspace:test/page:active").unwrap()
1706 }
1707
1708 fn provider_id() -> ToolProviderId {
1709 ToolProviderId::new("gate.browser-future").unwrap()
1710 }
1711
1712 fn gate_provider_id() -> ToolProviderId {
1713 ToolProviderId::new("gate.browser-shared").unwrap()
1714 }
1715
1716 fn capability_id() -> ToolCapabilityId {
1717 ToolCapabilityId::new("browser.page.snapshot").unwrap()
1718 }
1719
1720 fn provider() -> CapabilityProviderDescriptor {
1721 CapabilityProviderDescriptor {
1722 id: provider_id(),
1723 owner: CapabilityOwner::Consumer(consumer()),
1724 capabilities: vec![CapabilityDescriptor::new(
1725 capability_id(),
1726 CapabilityClass::Browser,
1727 "Return consumer-owned page state metadata",
1728 )
1729 .unwrap()],
1730 }
1731 }
1732
1733 fn gate_provider() -> CapabilityProviderDescriptor {
1734 CapabilityProviderDescriptor {
1735 id: gate_provider_id(),
1736 owner: CapabilityOwner::Gate,
1737 capabilities: vec![CapabilityDescriptor::new(
1738 capability_id(),
1739 CapabilityClass::Browser,
1740 "Return Gate-owned page state metadata",
1741 )
1742 .unwrap()],
1743 }
1744 }
1745
1746 fn grant(mode: GrantMode) -> PolicyGrant {
1747 PolicyGrant {
1748 key: PolicyKey {
1749 consumer_id: consumer(),
1750 actor_id: actor(),
1751 instance_id: instance(),
1752 generation: generation(1),
1753 provider_id: provider_id(),
1754 capability_id: capability_id(),
1755 resource_scope_id: resource_scope(),
1756 },
1757 mode,
1758 }
1759 }
1760
1761 fn scoped_grant(
1762 instance_id: AgentInstanceId,
1763 session_generation: SessionGeneration,
1764 resource: String,
1765 mode: GrantMode,
1766 ) -> PolicyGrant {
1767 PolicyGrant {
1768 key: PolicyKey {
1769 consumer_id: consumer(),
1770 actor_id: actor(),
1771 instance_id,
1772 generation: session_generation,
1773 provider_id: provider_id(),
1774 capability_id: capability_id(),
1775 resource_scope_id: ResourceScopeId::new(resource).unwrap(),
1776 },
1777 mode,
1778 }
1779 }
1780
1781 fn request(id: u64, generation: u64, deadline_tick: u64) -> CapabilityRequest {
1782 ConsumerBoundCapabilityRequest::new(
1783 consumer(),
1784 actor(),
1785 CapabilityRequestInput {
1786 local_id: CapabilityRequestId(id),
1787 instance_id: instance(),
1788 generation: SessionGeneration(generation),
1789 provider_id: provider_id(),
1790 capability_id: capability_id(),
1791 resource_scope_id: resource_scope(),
1792 approval_summary: "Read active page state".to_owned(),
1793 deadline_tick,
1794 payload: br#"{"scope":"active-page"}"#.to_vec(),
1795 },
1796 )
1797 }
1798
1799 fn request_key(id: u64) -> CapabilityRequestKey {
1800 CapabilityRequestKey {
1801 consumer_id: consumer(),
1802 actor_id: actor(),
1803 local_id: CapabilityRequestId(id),
1804 }
1805 }
1806
1807 fn request_for_client(
1808 id: u64,
1809 consumer_id: ConsumerId,
1810 actor_id: ToolActorId,
1811 ) -> CapabilityRequest {
1812 let mut request = request(id, 1, 100);
1813 request.consumer_id = consumer_id;
1814 request.actor_id = actor_id;
1815 request
1816 }
1817
1818 fn grant_for_client(
1819 consumer_id: ConsumerId,
1820 actor_id: ToolActorId,
1821 mode: GrantMode,
1822 ) -> PolicyGrant {
1823 let mut grant = grant(mode);
1824 grant.key.consumer_id = consumer_id;
1825 grant.key.actor_id = actor_id;
1826 grant
1827 }
1828
1829 fn accepted_sequence(engine: &ToolEngine, request_id: CapabilityRequestId) -> u64 {
1830 engine
1831 .requests
1832 .get(&request_key(request_id.0))
1833 .unwrap()
1834 .snapshot
1835 .accepted_sequence
1836 }
1837
1838 fn dummy_effect() -> CapabilityEffectEnvelope {
1839 CapabilityEffectEnvelope {
1840 sequence: 9_999,
1841 operation_id: ToolOperationId(9_999),
1842 request_key: request_key(9_999),
1843 instance_id: AgentInstanceId(9_999),
1844 generation: SessionGeneration(9_999),
1845 provider_id: provider_id(),
1846 deadline_tick: u64::MAX,
1847 effect: CapabilityEffect::Cancel {
1848 reason: InvocationCancelReason::DeadlineElapsed,
1849 },
1850 }
1851 }
1852
1853 fn fill_effect_queue(engine: &mut ToolEngine) {
1854 engine.effects.clear();
1855 engine.effects = vec![dummy_effect(); TOOL_EFFECTS_MAX];
1856 }
1857
1858 fn configured(mode: Option<GrantMode>) -> ToolEngine {
1859 let mut engine = ToolEngine::new();
1860 engine.register_provider(provider()).unwrap();
1861 engine.set_generation(instance(), generation(1)).unwrap();
1862 if let Some(mode) = mode {
1863 engine.set_grant(grant(mode)).unwrap();
1864 }
1865 engine
1866 }
1867
1868 struct FakeProvider;
1869
1870 impl FakeProvider {
1871 fn succeed(effect: &CapabilityEffectEnvelope) -> CapabilityObservationEnvelope {
1872 assert!(matches!(effect.effect, CapabilityEffect::Invoke { .. }));
1873 CapabilityObservationEnvelope {
1874 operation_id: effect.operation_id,
1875 request_key: effect.request_key.clone(),
1876 instance_id: effect.instance_id,
1877 generation: effect.generation,
1878 provider_id: effect.provider_id.clone(),
1879 observation: CapabilityObservation::Succeeded {
1880 result: CapabilityResult {
1881 metadata: CapabilityResultMetadata {
1882 byte_len: 2,
1883 media_type: Some("application/json".to_owned()),
1884 truncated: false,
1885 redacted_summary: Some("page snapshot captured".to_owned()),
1886 },
1887 delivery: CapabilityResultDelivery::Inline {
1888 bytes: b"{}".to_vec(),
1889 },
1890 },
1891 },
1892 }
1893 }
1894 }
1895
1896 #[test]
1897 fn provider_runtime_queries_preserve_registered_descriptors() {
1898 let mut engine = ToolEngine::new();
1899 assert!(!engine.provider_exists(&provider_id()));
1900 assert!(engine.provider_descriptor(&provider_id()).is_none());
1901 assert!(matches!(
1902 engine.detach_provider_runtime(&provider_id()),
1903 Err(ToolEngineError::UnknownRuntimeProvider { .. })
1904 ));
1905
1906 let descriptor = provider();
1907 engine.register_provider(descriptor.clone()).unwrap();
1908 assert!(engine.provider_exists(&provider_id()));
1909 assert_eq!(
1910 engine.provider_descriptor(&provider_id()),
1911 Some(&descriptor)
1912 );
1913 assert_eq!(engine.detach_provider_runtime(&provider_id()).unwrap(), 0);
1914 assert_eq!(
1915 engine.provider_descriptor(&provider_id()),
1916 Some(&descriptor)
1917 );
1918 }
1919
1920 #[test]
1921 fn provider_detach_closes_pre_dispatch_request_without_removing_policy() {
1922 let mut engine = configured(Some(GrantMode::RequireApproval));
1923 let descriptor = engine.provider_descriptor(&provider_id()).unwrap().clone();
1924 let policy = grant(GrantMode::RequireApproval).key;
1925 engine.request(request(1, 1, 100)).unwrap();
1926 assert!(!engine.requests[&request_key(1)].request.payload.is_empty());
1927
1928 assert_eq!(engine.detach_provider_runtime(&provider_id()).unwrap(), 1);
1929 assert!(engine.requests[&request_key(1)].request.payload.is_empty());
1930 assert!(matches!(
1931 engine.requests[&request_key(1)].snapshot.status,
1932 CapabilityRequestStatus::ProviderDetached {
1933 operation_id: None,
1934 cancellation: CancellationDisposition::NotRequired,
1935 }
1936 ));
1937 assert_eq!(
1938 engine.provider_descriptor(&provider_id()),
1939 Some(&descriptor)
1940 );
1941 assert_eq!(
1942 engine.grants.get(&policy),
1943 Some(&GrantMode::RequireApproval)
1944 );
1945 assert!(engine.drain_effects().is_empty());
1946 assert!(matches!(
1947 engine.drain_completions().completions[0].outcome,
1948 CapabilityTerminalOutcome::ProviderDetached {
1949 cancellation: CancellationDisposition::NotRequired,
1950 }
1951 ));
1952 assert!(engine.snapshot().audit_events.iter().any(|event| matches!(
1953 event.event,
1954 ToolAuditEventKind::RequestProviderDetached {
1955 operation_id: None,
1956 cancellation: CancellationDisposition::NotRequired,
1957 }
1958 )));
1959 }
1960
1961 #[test]
1962 fn provider_detach_removes_all_queued_invokes_and_fences_late_results() {
1963 let mut engine = configured(Some(GrantMode::Allow));
1964 engine.request(request(1, 1, 100)).unwrap();
1965 engine.request(request(2, 1, 100)).unwrap();
1966 let late_invoke = engine.effects[0].clone();
1967
1968 assert_eq!(engine.detach_provider_runtime(&provider_id()).unwrap(), 2);
1969 assert!(engine.drain_effects().is_empty());
1970 for request_id in [1, 2] {
1971 assert!(matches!(
1972 engine.requests[&request_key(request_id)].snapshot.status,
1973 CapabilityRequestStatus::ProviderDetached {
1974 operation_id: Some(_),
1975 cancellation: CancellationDisposition::QueuedInvokeRemoved,
1976 }
1977 ));
1978 }
1979 let completion_batch = engine.drain_completions();
1980 assert_eq!(completion_batch.completions.len(), 2);
1981 assert!(completion_batch
1982 .completions
1983 .iter()
1984 .all(|completion| matches!(
1985 completion.outcome,
1986 CapabilityTerminalOutcome::ProviderDetached {
1987 cancellation: CancellationDisposition::QueuedInvokeRemoved,
1988 }
1989 )));
1990
1991 engine
1992 .apply_observation(FakeProvider::succeed(&late_invoke))
1993 .unwrap();
1994 assert!(matches!(
1995 engine.requests[&request_key(1)].snapshot.status,
1996 CapabilityRequestStatus::ProviderDetached {
1997 cancellation: CancellationDisposition::QueuedInvokeRemoved,
1998 ..
1999 }
2000 ));
2001 assert!(matches!(
2002 engine.snapshot().audit_events.last().unwrap().event,
2003 ToolAuditEventKind::ObservationIgnored {
2004 reason: ObservationIgnoredReason::RequestNotDispatched,
2005 ..
2006 }
2007 ));
2008 }
2009
2010 #[test]
2011 fn provider_detach_marks_drained_execution_unconfirmed_without_cancel_effect() {
2012 let mut engine = configured(Some(GrantMode::Allow));
2013 engine.request(request(1, 1, 100)).unwrap();
2014 let invoke = engine.drain_effects().pop().unwrap();
2015
2016 assert_eq!(engine.detach_provider_runtime(&provider_id()).unwrap(), 1);
2017 assert!(engine.drain_effects().is_empty());
2018 assert!(matches!(
2019 engine.requests[&request_key(1)].snapshot.status,
2020 CapabilityRequestStatus::ProviderDetached {
2021 operation_id: Some(operation_id),
2022 cancellation: CancellationDisposition::ProviderDetachedUnconfirmed,
2023 } if operation_id == invoke.operation_id
2024 ));
2025 assert!(matches!(
2026 engine.drain_completions().completions[0].outcome,
2027 CapabilityTerminalOutcome::ProviderDetached {
2028 cancellation: CancellationDisposition::ProviderDetachedUnconfirmed,
2029 }
2030 ));
2031 }
2032
2033 #[test]
2034 fn provider_detach_purges_cancel_already_queued_for_terminal_request() {
2035 let mut engine = configured(Some(GrantMode::Allow));
2036 let policy_key = grant(GrantMode::Allow).key;
2037 engine.request(request(1, 1, 100)).unwrap();
2038 let invoke = engine.drain_effects().pop().unwrap();
2039
2040 assert!(engine.revoke_grant(&policy_key).unwrap());
2041 assert!(matches!(
2042 &engine.effects[0].effect,
2043 CapabilityEffect::Cancel {
2044 reason: InvocationCancelReason::GrantRevoked
2045 }
2046 ));
2047 assert_eq!(engine.effects[0].operation_id, invoke.operation_id);
2048
2049 assert_eq!(engine.detach_provider_runtime(&provider_id()).unwrap(), 0);
2050 assert!(engine.drain_effects().is_empty());
2051 assert!(matches!(
2052 engine.requests[&request_key(1)].snapshot.status,
2053 CapabilityRequestStatus::GrantRevoked { .. }
2054 ));
2055 }
2056
2057 #[test]
2058 fn policy_is_deny_by_default_and_releases_no_effect() {
2059 let mut engine = configured(None);
2060 assert_eq!(
2061 engine.request(request(1, 1, 10)).unwrap(),
2062 PolicyDecision::Deny(PolicyDenial::MissingGrant)
2063 );
2064 assert!(engine.drain_effects().is_empty());
2065 assert!(matches!(
2066 engine.drain_completions().completions[0].outcome,
2067 CapabilityTerminalOutcome::PolicyDenied { .. }
2068 ));
2069 let snapshot = engine.snapshot();
2070 assert_eq!(snapshot.requests[0].payload_bytes, 23);
2071 assert!(matches!(
2072 snapshot.requests[0].status,
2073 CapabilityRequestStatus::Denied {
2074 reason: PolicyDenial::MissingGrant
2075 }
2076 ));
2077 }
2078
2079 #[test]
2080 fn approval_releases_exactly_one_typed_effect() {
2081 let mut engine = configured(Some(GrantMode::RequireApproval));
2082 assert_eq!(
2083 engine.request(request(1, 1, 10)).unwrap(),
2084 PolicyDecision::RequireApproval
2085 );
2086 assert!(engine.drain_effects().is_empty());
2087 let accepted_sequence = accepted_sequence(&engine, CapabilityRequestId(1));
2088 engine
2089 .resolve_approval(ApprovalResolution {
2090 request_key: request_key(1),
2091 accepted_sequence,
2092 instance_id: instance(),
2093 generation: generation(1),
2094 decision: ApprovalDecision::ApproveOnce,
2095 })
2096 .unwrap();
2097 let effects = engine.drain_effects();
2098 assert_eq!(effects.len(), 1);
2099 assert!(matches!(effects[0].effect, CapabilityEffect::Invoke { .. }));
2100 assert!(engine
2101 .requests
2102 .get(&request_key(1))
2103 .unwrap()
2104 .request
2105 .payload
2106 .is_empty());
2107 }
2108
2109 #[test]
2110 fn grant_is_exact_and_does_not_authorize_another_actor() {
2111 let mut engine = configured(Some(GrantMode::Allow));
2112 let mut ungranted = request(1, 1, 10);
2113 ungranted.actor_id = ToolActorId::new("consumer.other-agent").unwrap();
2114 assert_eq!(
2115 engine.request(ungranted).unwrap(),
2116 PolicyDecision::Deny(PolicyDenial::MissingGrant)
2117 );
2118 assert!(engine.drain_effects().is_empty());
2119 }
2120
2121 #[test]
2122 fn grant_is_exact_across_consumer_resource_and_session_generation() {
2123 let mut engine = configured(Some(GrantMode::Allow));
2124 let mut other_consumer = request(1, 1, 10);
2125 other_consumer.consumer_id = ConsumerId::new("station.other").unwrap();
2126 assert_eq!(
2127 engine.request(other_consumer).unwrap(),
2128 PolicyDecision::Deny(PolicyDenial::ProviderOwnerMismatch)
2129 );
2130 let mut other_resource = request(2, 1, 10);
2131 other_resource.request.resource_scope_id =
2132 ResourceScopeId::new("workspace:test/page:other").unwrap();
2133 assert_eq!(
2134 engine.request(other_resource).unwrap(),
2135 PolicyDecision::Deny(PolicyDenial::MissingGrant)
2136 );
2137 engine.set_generation(instance(), generation(2)).unwrap();
2138 assert_eq!(
2139 engine.request(request(3, 2, 10)).unwrap(),
2140 PolicyDecision::Deny(PolicyDenial::MissingGrant)
2141 );
2142 assert!(engine.drain_effects().is_empty());
2143 }
2144
2145 #[test]
2146 fn consumer_owned_provider_rejects_owner_mismatch_at_grant_and_request() {
2147 let mut engine = configured(None);
2148 let mut mismatched_grant = grant(GrantMode::Allow);
2149 mismatched_grant.key.consumer_id = ConsumerId::new("station.other").unwrap();
2150 assert!(matches!(
2151 engine.set_grant(mismatched_grant),
2152 Err(ToolEngineError::ProviderOwnerMismatch { .. })
2153 ));
2154
2155 let mut mismatched_request = request(1, 1, 10);
2156 mismatched_request.consumer_id = ConsumerId::new("station.other").unwrap();
2157 assert_eq!(
2158 engine.request(mismatched_request).unwrap(),
2159 PolicyDecision::Deny(PolicyDenial::ProviderOwnerMismatch)
2160 );
2161 assert!(engine.drain_effects().is_empty());
2162 }
2163
2164 #[test]
2165 fn gate_owned_provider_is_shareable_only_by_exact_grant() {
2166 let mut engine = ToolEngine::new();
2167 engine.register_provider(gate_provider()).unwrap();
2168 engine.set_generation(instance(), generation(1)).unwrap();
2169 let mut exact_grant = grant(GrantMode::Allow);
2170 exact_grant.key.provider_id = gate_provider_id();
2171 engine.set_grant(exact_grant).unwrap();
2172
2173 let mut exact_request = request(1, 1, 10);
2174 exact_request.request.provider_id = gate_provider_id();
2175 assert_eq!(
2176 engine.request(exact_request).unwrap(),
2177 PolicyDecision::Allow
2178 );
2179 engine.drain_effects();
2180 let mut other_consumer = request(2, 1, 10);
2181 other_consumer.request.provider_id = gate_provider_id();
2182 other_consumer.consumer_id = ConsumerId::new("station.other").unwrap();
2183 assert_eq!(
2184 engine.request(other_consumer).unwrap(),
2185 PolicyDecision::Deny(PolicyDenial::MissingGrant)
2186 );
2187 assert!(engine.drain_effects().is_empty());
2188 }
2189
2190 #[test]
2191 fn grant_replacement_revokes_open_approval_and_queued_invoke() {
2192 let mut approval_engine = configured(Some(GrantMode::RequireApproval));
2193 approval_engine.request(request(1, 1, 10)).unwrap();
2194 approval_engine.set_grant(grant(GrantMode::Allow)).unwrap();
2195 assert!(matches!(
2196 approval_engine.snapshot().requests[0].status,
2197 CapabilityRequestStatus::GrantRevoked {
2198 cancellation: CancellationDisposition::NotRequired,
2199 ..
2200 }
2201 ));
2202 assert!(approval_engine.drain_effects().is_empty());
2203
2204 let mut invoke_engine = configured(Some(GrantMode::Allow));
2205 invoke_engine.request(request(1, 1, 10)).unwrap();
2206 invoke_engine
2207 .set_grant(grant(GrantMode::RequireApproval))
2208 .unwrap();
2209 assert!(matches!(
2210 invoke_engine.snapshot().requests[0].status,
2211 CapabilityRequestStatus::GrantRevoked {
2212 cancellation: CancellationDisposition::QueuedInvokeRemoved,
2213 ..
2214 }
2215 ));
2216 assert!(invoke_engine.drain_effects().is_empty());
2217 }
2218
2219 #[test]
2220 fn revoking_grant_closes_pending_approval_and_erases_payload() {
2221 let mut engine = configured(Some(GrantMode::RequireApproval));
2222 engine.request(request(1, 1, 10)).unwrap();
2223 assert!(!engine
2224 .requests
2225 .get(&request_key(1))
2226 .unwrap()
2227 .request
2228 .payload
2229 .is_empty());
2230 assert!(engine.revoke_grant(&grant(GrantMode::Allow).key).unwrap());
2231 assert!(engine
2232 .requests
2233 .get(&request_key(1))
2234 .unwrap()
2235 .request
2236 .payload
2237 .is_empty());
2238 assert!(matches!(
2239 engine.snapshot().requests[0].status,
2240 CapabilityRequestStatus::GrantRevoked {
2241 operation_id: None,
2242 cancellation: CancellationDisposition::NotRequired,
2243 }
2244 ));
2245 assert!(matches!(
2246 engine.resolve_approval(ApprovalResolution {
2247 request_key: request_key(1),
2248 accepted_sequence: accepted_sequence(&engine, CapabilityRequestId(1)),
2249 instance_id: instance(),
2250 generation: generation(1),
2251 decision: ApprovalDecision::ApproveOnce,
2252 }),
2253 Err(ToolEngineError::RequestNotAwaitingApproval { .. })
2254 ));
2255 assert!(engine.drain_effects().is_empty());
2256 assert!(matches!(
2257 engine.drain_completions().completions[0].outcome,
2258 CapabilityTerminalOutcome::GrantRevoked { .. }
2259 ));
2260 }
2261
2262 #[test]
2263 fn fake_provider_success_closes_only_the_matching_operation() {
2264 let mut engine = configured(Some(GrantMode::Allow));
2265 engine.request(request(1, 1, 10)).unwrap();
2266 let effect = engine.drain_effects().pop().unwrap();
2267 let mut mismatched = FakeProvider::succeed(&effect);
2268 mismatched.operation_id = ToolOperationId(effect.operation_id.0 + 1);
2269 engine.apply_observation(mismatched).unwrap();
2270 assert!(matches!(
2271 engine.snapshot().requests[0].status,
2272 CapabilityRequestStatus::Dispatched { .. }
2273 ));
2274 engine
2275 .apply_observation(FakeProvider::succeed(&effect))
2276 .unwrap();
2277 assert!(matches!(
2278 engine.snapshot().requests[0].status,
2279 CapabilityRequestStatus::Succeeded { .. }
2280 ));
2281 let completion = engine.drain_completions().completions.pop().unwrap();
2282 assert_eq!(completion.operation_id, Some(effect.operation_id));
2283 assert!(matches!(
2284 completion.outcome,
2285 CapabilityTerminalOutcome::Succeeded {
2286 result: CapabilityResult {
2287 delivery: CapabilityResultDelivery::Inline { ref bytes },
2288 ..
2289 }
2290 } if bytes == b"{}"
2291 ));
2292 }
2293
2294 #[test]
2295 fn generation_advance_cancels_and_stale_success_cannot_resurrect_request() {
2296 let mut engine = configured(Some(GrantMode::Allow));
2297 engine.request(request(1, 1, 10)).unwrap();
2298 let invoke = engine.drain_effects().pop().unwrap();
2299 engine.set_generation(instance(), generation(2)).unwrap();
2300 let cancel = engine.drain_effects().pop().unwrap();
2301 assert_eq!(cancel.operation_id, invoke.operation_id);
2302 assert!(matches!(
2303 cancel.effect,
2304 CapabilityEffect::Cancel {
2305 reason: InvocationCancelReason::GenerationSuperseded
2306 }
2307 ));
2308 engine
2309 .apply_observation(FakeProvider::succeed(&invoke))
2310 .unwrap();
2311 let snapshot = engine.snapshot();
2312 assert!(matches!(
2313 snapshot.requests[0].status,
2314 CapabilityRequestStatus::Superseded {
2315 current_generation: SessionGeneration(2),
2316 ..
2317 }
2318 ));
2319 assert!(snapshot.audit_events.iter().any(|event| matches!(
2320 event.event,
2321 ToolAuditEventKind::ObservationIgnored {
2322 reason: ObservationIgnoredReason::StaleGeneration,
2323 ..
2324 }
2325 )));
2326 assert!(matches!(
2327 engine.drain_completions().completions[0].outcome,
2328 CapabilityTerminalOutcome::Superseded { .. }
2329 ));
2330 }
2331
2332 #[test]
2333 fn deadline_cancels_dispatched_work_and_late_result_is_ignored() {
2334 let mut engine = configured(Some(GrantMode::Allow));
2335 engine.request(request(1, 1, 5)).unwrap();
2336 let invoke = engine.drain_effects().pop().unwrap();
2337 engine.advance_time(5).unwrap();
2338 let cancel = engine.drain_effects().pop().unwrap();
2339 assert_eq!(cancel.operation_id, invoke.operation_id);
2340 assert!(matches!(
2341 cancel.effect,
2342 CapabilityEffect::Cancel {
2343 reason: InvocationCancelReason::DeadlineElapsed
2344 }
2345 ));
2346 engine
2347 .apply_observation(FakeProvider::succeed(&invoke))
2348 .unwrap();
2349 assert!(matches!(
2350 engine.snapshot().requests[0].status,
2351 CapabilityRequestStatus::TimedOut { .. }
2352 ));
2353 assert!(matches!(
2354 engine.drain_completions().completions[0].outcome,
2355 CapabilityTerminalOutcome::TimedOut { .. }
2356 ));
2357 }
2358
2359 #[test]
2360 fn oversized_input_is_rejected_before_policy_or_effect() {
2361 let mut engine = configured(Some(GrantMode::Allow));
2362 let mut oversized = request(1, 1, 10);
2363 oversized.request.payload = vec![0; TOOL_PAYLOAD_MAX_BYTES + 1];
2364 assert!(matches!(
2365 engine.request(oversized),
2366 Err(ToolEngineError::Validation(ToolValidationError::TooLarge {
2367 field: "tool request payload",
2368 ..
2369 }))
2370 ));
2371 assert!(engine.snapshot().requests.is_empty());
2372 assert!(engine.drain_effects().is_empty());
2373 }
2374
2375 #[test]
2376 fn wire_deserialization_cannot_bypass_bounded_identifier_constructor() {
2377 let invalid = format!("\"{}\"", "x".repeat(crate::TOOL_ACTOR_ID_MAX_BYTES + 1));
2378 assert!(serde_json::from_str::<ToolActorId>(&invalid).is_err());
2379 assert!(serde_json::from_str::<ToolActorId>("\"contains space\"").is_err());
2380 }
2381
2382 #[test]
2383 fn approval_nonce_rejects_aba_after_terminal_eviction_and_id_reuse() {
2384 let mut engine = configured(Some(GrantMode::RequireApproval));
2385 engine.request(request(1, 1, 100)).unwrap();
2386 let old_sequence = accepted_sequence(&engine, CapabilityRequestId(1));
2387 engine
2388 .resolve_approval(ApprovalResolution {
2389 request_key: request_key(1),
2390 accepted_sequence: old_sequence,
2391 instance_id: instance(),
2392 generation: generation(1),
2393 decision: ApprovalDecision::Deny,
2394 })
2395 .unwrap();
2396 engine.revoke_grant(&grant(GrantMode::Allow).key).unwrap();
2397 for id in 2..=TOOL_REQUESTS_MAX as u64 {
2398 engine.request(request(id, 1, 100)).unwrap();
2399 }
2400 engine.set_grant(grant(GrantMode::RequireApproval)).unwrap();
2401 engine
2402 .request(request(TOOL_REQUESTS_MAX as u64 + 1, 1, 100))
2403 .unwrap();
2404 engine.request(request(1, 1, 100)).unwrap();
2405 let new_sequence = accepted_sequence(&engine, CapabilityRequestId(1));
2406 assert_ne!(old_sequence, new_sequence);
2407 assert!(matches!(
2408 engine.resolve_approval(ApprovalResolution {
2409 request_key: request_key(1),
2410 accepted_sequence: old_sequence,
2411 instance_id: instance(),
2412 generation: generation(1),
2413 decision: ApprovalDecision::ApproveOnce,
2414 }),
2415 Err(ToolEngineError::ApprovalNonceMismatch {
2416 expected,
2417 actual,
2418 ..
2419 }) if expected == new_sequence && actual == old_sequence
2420 ));
2421 assert!(engine.drain_effects().is_empty());
2422 }
2423
2424 #[test]
2425 fn safety_transitions_close_authority_even_when_effect_queue_is_full() {
2426 let mut generation_engine = configured(Some(GrantMode::Allow));
2427 generation_engine.request(request(1, 1, 10)).unwrap();
2428 generation_engine.drain_effects();
2429 fill_effect_queue(&mut generation_engine);
2430 generation_engine
2431 .set_generation(instance(), generation(2))
2432 .unwrap();
2433 assert_eq!(generation_engine.snapshot().generations[0].1, generation(2));
2434 assert!(matches!(
2435 generation_engine.snapshot().requests[0].status,
2436 CapabilityRequestStatus::Superseded {
2437 cancellation: CancellationDisposition::DroppedQueueFull,
2438 ..
2439 }
2440 ));
2441
2442 let mut revoke_engine = configured(Some(GrantMode::Allow));
2443 revoke_engine.request(request(1, 1, 10)).unwrap();
2444 revoke_engine.drain_effects();
2445 fill_effect_queue(&mut revoke_engine);
2446 assert!(revoke_engine
2447 .revoke_grant(&grant(GrantMode::Allow).key)
2448 .unwrap());
2449 assert!(revoke_engine.snapshot().grants.is_empty());
2450 assert!(matches!(
2451 revoke_engine.snapshot().requests[0].status,
2452 CapabilityRequestStatus::GrantRevoked {
2453 cancellation: CancellationDisposition::DroppedQueueFull,
2454 ..
2455 }
2456 ));
2457
2458 let mut time_engine = configured(Some(GrantMode::Allow));
2459 time_engine.request(request(1, 1, 5)).unwrap();
2460 while time_engine.effects.len() < TOOL_EFFECTS_MAX {
2461 time_engine.effects.push(dummy_effect());
2462 }
2463 time_engine.advance_time(5).unwrap();
2464 assert_eq!(time_engine.snapshot().current_tick, 5);
2465 assert!(matches!(
2466 time_engine.snapshot().requests[0].status,
2467 CapabilityRequestStatus::TimedOut {
2468 cancellation: CancellationDisposition::QueuedInvokeRemoved,
2469 ..
2470 }
2471 ));
2472 assert_eq!(time_engine.effects.len(), TOOL_EFFECTS_MAX - 1);
2473 }
2474
2475 #[test]
2476 fn generation_advance_purges_only_stale_instance_grants_and_reuses_capacity() {
2477 let mut engine = configured(None);
2478 let other_instance = AgentInstanceId(42);
2479 engine
2480 .set_generation(other_instance, generation(1))
2481 .unwrap();
2482 let other_grant = scoped_grant(
2483 other_instance,
2484 generation(1),
2485 "workspace:other/current".to_owned(),
2486 GrantMode::Allow,
2487 );
2488 engine.set_grant(other_grant.clone()).unwrap();
2489
2490 let stale_count = TOOL_POLICIES_MAX - 2;
2491 for index in 0..stale_count {
2492 engine
2493 .set_grant(scoped_grant(
2494 instance(),
2495 generation(1),
2496 format!("workspace:test/stale:{index}"),
2497 GrantMode::Allow,
2498 ))
2499 .unwrap();
2500 }
2501 let current_generation_grant = scoped_grant(
2502 instance(),
2503 generation(2),
2504 "workspace:test/current".to_owned(),
2505 GrantMode::Allow,
2506 );
2507 engine.grants.insert(
2508 current_generation_grant.key.clone(),
2509 current_generation_grant.mode,
2510 );
2511 assert_eq!(engine.grants.len(), TOOL_POLICIES_MAX);
2512
2513 engine.set_generation(instance(), generation(2)).unwrap();
2514 assert_eq!(engine.grants.len(), 2);
2515 assert!(engine.grants.contains_key(&other_grant.key));
2516 assert!(engine.grants.contains_key(¤t_generation_grant.key));
2517 assert!(engine.snapshot().audit_events.iter().any(|event| matches!(
2518 event.event,
2519 ToolAuditEventKind::GenerationAdvanced {
2520 instance_id,
2521 current: SessionGeneration(2),
2522 purged_grant_count,
2523 ..
2524 } if instance_id == instance() && purged_grant_count == stale_count
2525 )));
2526 engine
2527 .set_grant(scoped_grant(
2528 instance(),
2529 generation(2),
2530 "workspace:test/reused-capacity".to_owned(),
2531 GrantMode::Allow,
2532 ))
2533 .unwrap();
2534 assert_eq!(engine.grants.len(), 3);
2535 }
2536
2537 #[test]
2538 fn failed_effect_preflight_does_not_evict_terminal_request() {
2539 let mut engine = configured(None);
2540 for id in 1..=TOOL_REQUESTS_MAX as u64 {
2541 engine.request(request(id, 1, 100)).unwrap();
2542 }
2543 engine.set_grant(grant(GrantMode::Allow)).unwrap();
2544 fill_effect_queue(&mut engine);
2545 assert!(matches!(
2546 engine.request(request(TOOL_REQUESTS_MAX as u64 + 1, 1, 100)),
2547 Err(ToolEngineError::EffectCapacityExceeded)
2548 ));
2549 assert_eq!(engine.snapshot().requests.len(), TOOL_REQUESTS_MAX);
2550 assert!(engine.requests.contains_key(&request_key(1)));
2551 }
2552
2553 #[test]
2554 fn debug_redacts_raw_payload_and_control_characters_are_rejected() {
2555 let secret_payload = vec![13, 37, 201, 222, 173, 190, 239];
2556 let secret_payload_debug = format!("{secret_payload:?}");
2557 let mut raw_request = request(1, 1, 10);
2558 raw_request.request.payload = secret_payload.clone();
2559 assert!(!format!("{raw_request:?}").contains(&secret_payload_debug));
2560 let raw_effect = CapabilityEffect::Invoke {
2561 consumer_id: consumer(),
2562 actor_id: actor(),
2563 capability_id: capability_id(),
2564 resource_scope_id: resource_scope(),
2565 payload: secret_payload.clone(),
2566 };
2567 assert!(!format!("{raw_effect:?}").contains(&secret_payload_debug));
2568 let raw_effect_envelope = CapabilityEffectEnvelope {
2569 sequence: 1,
2570 operation_id: ToolOperationId(1),
2571 request_key: request_key(1),
2572 instance_id: instance(),
2573 generation: generation(1),
2574 provider_id: provider_id(),
2575 deadline_tick: 10,
2576 effect: raw_effect,
2577 };
2578 assert!(!format!("{raw_effect_envelope:?}").contains(&secret_payload_debug));
2579
2580 let secret_inline = vec![91, 17, 233, 44, 155];
2581 let secret_inline_debug = format!("{secret_inline:?}");
2582 let inline_delivery = CapabilityResultDelivery::Inline {
2583 bytes: secret_inline.clone(),
2584 };
2585 let inline_result = CapabilityResult {
2586 metadata: CapabilityResultMetadata {
2587 byte_len: secret_inline.len() as u64,
2588 media_type: None,
2589 truncated: false,
2590 redacted_summary: None,
2591 },
2592 delivery: inline_delivery.clone(),
2593 };
2594 let inline_completion = CapabilityCompletionEnvelope {
2595 sequence: 1,
2596 accepted_sequence: 1,
2597 operation_id: Some(ToolOperationId(1)),
2598 request_key: request_key(1),
2599 instance_id: instance(),
2600 generation: generation(1),
2601 provider_id: provider_id(),
2602 outcome: CapabilityTerminalOutcome::Succeeded {
2603 result: inline_result.clone(),
2604 },
2605 };
2606 let inline_observation = CapabilityObservation::Succeeded {
2607 result: inline_result.clone(),
2608 };
2609 let inline_observation_envelope = CapabilityObservationEnvelope {
2610 operation_id: ToolOperationId(1),
2611 request_key: request_key(1),
2612 instance_id: instance(),
2613 generation: generation(1),
2614 provider_id: provider_id(),
2615 observation: inline_observation.clone(),
2616 };
2617 for rendered in [
2618 format!("{inline_delivery:?}"),
2619 format!("{inline_result:?}"),
2620 format!("{inline_completion:?}"),
2621 format!("{inline_observation:?}"),
2622 format!("{inline_observation_envelope:?}"),
2623 ] {
2624 assert!(!rendered.contains(&secret_inline_debug));
2625 }
2626
2627 let secret_reference = "opaque://DO_NOT_FORMAT_REFERENCE";
2628 let reference_delivery = CapabilityResultDelivery::OpaqueReference {
2629 reference: secret_reference.to_owned(),
2630 };
2631 let reference_result = CapabilityResult {
2632 metadata: CapabilityResultMetadata {
2633 byte_len: 128,
2634 media_type: None,
2635 truncated: false,
2636 redacted_summary: None,
2637 },
2638 delivery: reference_delivery.clone(),
2639 };
2640 let reference_completion = CapabilityCompletionEnvelope {
2641 sequence: 2,
2642 accepted_sequence: 2,
2643 operation_id: Some(ToolOperationId(2)),
2644 request_key: request_key(2),
2645 instance_id: instance(),
2646 generation: generation(1),
2647 provider_id: provider_id(),
2648 outcome: CapabilityTerminalOutcome::Succeeded {
2649 result: reference_result.clone(),
2650 },
2651 };
2652 let reference_observation = CapabilityObservation::Succeeded {
2653 result: reference_result.clone(),
2654 };
2655 let reference_observation_envelope = CapabilityObservationEnvelope {
2656 operation_id: ToolOperationId(2),
2657 request_key: request_key(2),
2658 instance_id: instance(),
2659 generation: generation(1),
2660 provider_id: provider_id(),
2661 observation: reference_observation.clone(),
2662 };
2663 for rendered in [
2664 format!("{reference_delivery:?}"),
2665 format!("{reference_result:?}"),
2666 format!("{reference_completion:?}"),
2667 format!("{reference_observation:?}"),
2668 format!("{reference_observation_envelope:?}"),
2669 ] {
2670 assert!(!rendered.contains(secret_reference));
2671 }
2672
2673 let mut engine = configured(Some(GrantMode::RequireApproval));
2674 engine.request(raw_request).unwrap();
2675 assert!(!format!("{engine:?}").contains(&secret_payload_debug));
2676
2677 let mut unsafe_summary = request(2, 1, 10);
2678 unsafe_summary.request.approval_summary = "read page\u{1b}[31m".to_owned();
2679 assert!(matches!(
2680 engine.request(unsafe_summary),
2681 Err(ToolEngineError::Validation(
2682 ToolValidationError::ControlCharacter {
2683 field: "tool approval summary"
2684 }
2685 ))
2686 ));
2687 let mut whitespace_summary = request(3, 1, 10);
2688 whitespace_summary.request.approval_summary = " \t ".to_owned();
2689 assert!(matches!(
2690 engine.request(whitespace_summary),
2691 Err(ToolEngineError::Validation(ToolValidationError::Required {
2692 field: "tool approval summary"
2693 }))
2694 ));
2695 assert!(matches!(
2696 CapabilityResultMetadata {
2697 byte_len: 0,
2698 media_type: None,
2699 truncated: false,
2700 redacted_summary: Some("unsafe\u{1b}".to_owned()),
2701 }
2702 .validate(),
2703 Err(ToolValidationError::ControlCharacter {
2704 field: "tool result redacted summary"
2705 })
2706 ));
2707 assert!(matches!(
2708 ToolFailure {
2709 kind: crate::ToolFailureKind::Execution,
2710 redacted_message: Some("unsafe\nmessage".to_owned()),
2711 }
2712 .validate(),
2713 Err(ToolValidationError::ControlCharacter {
2714 field: "tool failure redacted message"
2715 })
2716 ));
2717 }
2718
2719 #[test]
2720 fn capability_admission_rejects_shell_filesystem_and_mcp_namespaces() {
2721 for id in [
2722 "shell.exec",
2723 "filesystem.read",
2724 "mcp.call",
2725 "browser.mcp.call",
2726 "browser.filesystem-read",
2727 ] {
2728 assert!(matches!(
2729 CapabilityDescriptor::new(
2730 ToolCapabilityId::new(id).unwrap(),
2731 CapabilityClass::Browser,
2732 "unsafe capability",
2733 ),
2734 Err(ToolValidationError::CapabilityOutsideAdmission { .. })
2735 ));
2736 }
2737 assert!(CapabilityDescriptor::new(
2738 ToolCapabilityId::new("consumer-state.selection.read").unwrap(),
2739 CapabilityClass::ConsumerState,
2740 "Read bounded consumer state",
2741 )
2742 .is_ok());
2743 }
2744
2745 #[test]
2746 fn replay_is_deterministic_and_audit_never_contains_payload() {
2747 fn replay() -> (ToolEngineSnapshot, Vec<CapabilityEffectEnvelope>) {
2748 let mut engine = configured(Some(GrantMode::Allow));
2749 engine.request(request(1, 1, 10)).unwrap();
2750 let effect = engine.drain_effects().pop().unwrap();
2751 engine
2752 .apply_observation(FakeProvider::succeed(&effect))
2753 .unwrap();
2754 (engine.snapshot(), vec![effect])
2755 }
2756
2757 let first = replay();
2758 let second = replay();
2759 assert_eq!(first, second);
2760 assert_eq!(first.0.requests[0].payload_bytes, 23);
2761 }
2762
2763 #[test]
2764 fn provider_failure_and_approval_denial_emit_terminal_completions() {
2765 let mut failed = configured(Some(GrantMode::Allow));
2766 failed.request(request(1, 1, 100)).unwrap();
2767 let effect = failed.drain_effects().pop().unwrap();
2768 failed
2769 .apply_observation(CapabilityObservationEnvelope {
2770 operation_id: effect.operation_id,
2771 request_key: effect.request_key,
2772 instance_id: effect.instance_id,
2773 generation: effect.generation,
2774 provider_id: effect.provider_id,
2775 observation: CapabilityObservation::Failed {
2776 failure: ToolFailure {
2777 kind: ToolFailureKind::Execution,
2778 redacted_message: Some("provider failed".to_owned()),
2779 },
2780 },
2781 })
2782 .unwrap();
2783 assert!(matches!(
2784 failed.drain_completions().completions[0].outcome,
2785 CapabilityTerminalOutcome::Failed { .. }
2786 ));
2787
2788 let mut denied = configured(Some(GrantMode::RequireApproval));
2789 denied.request(request(1, 1, 100)).unwrap();
2790 denied
2791 .resolve_approval(ApprovalResolution {
2792 request_key: request_key(1),
2793 accepted_sequence: accepted_sequence(&denied, CapabilityRequestId(1)),
2794 instance_id: instance(),
2795 generation: generation(1),
2796 decision: ApprovalDecision::Deny,
2797 })
2798 .unwrap();
2799 assert!(matches!(
2800 denied.drain_completions().completions[0].outcome,
2801 CapabilityTerminalOutcome::ApprovalDenied
2802 ));
2803 }
2804
2805 #[test]
2806 fn same_local_id_is_scoped_by_consumer_and_actor() {
2807 let consumer_b = ConsumerId::new("station.other").unwrap();
2808 let actor_b = ToolActorId::new("consumer.other-agent").unwrap();
2809 let mut engine = ToolEngine::new();
2810 engine.register_provider(gate_provider()).unwrap();
2811 engine.set_generation(instance(), generation(1)).unwrap();
2812
2813 let mut grant_a = grant_for_client(consumer(), actor(), GrantMode::Allow);
2814 grant_a.key.provider_id = gate_provider_id();
2815 let mut grant_b = grant_for_client(consumer_b.clone(), actor_b.clone(), GrantMode::Allow);
2816 grant_b.key.provider_id = gate_provider_id();
2817 engine.set_grant(grant_a).unwrap();
2818 engine.set_grant(grant_b).unwrap();
2819
2820 let mut request_a = request_for_client(7, consumer(), actor());
2821 request_a.request.provider_id = gate_provider_id();
2822 let mut request_b = request_for_client(7, consumer_b.clone(), actor_b.clone());
2823 request_b.request.provider_id = gate_provider_id();
2824 assert_eq!(engine.request(request_a).unwrap(), PolicyDecision::Allow);
2825 assert_eq!(engine.request(request_b).unwrap(), PolicyDecision::Allow);
2826 let effects = engine.drain_effects();
2827 assert_eq!(effects.len(), 2);
2828 assert_ne!(effects[0].request_key, effects[1].request_key);
2829 assert_eq!(
2830 effects[0].request_key.local_id,
2831 effects[1].request_key.local_id
2832 );
2833 }
2834
2835 #[test]
2836 fn approval_target_is_exact_and_forgery_does_not_mutate() {
2837 let mut engine = configured(Some(GrantMode::RequireApproval));
2838 engine.request(request(1, 1, 10)).unwrap();
2839 let before = engine.snapshot();
2840 let forged_key = CapabilityRequestKey {
2841 consumer_id: consumer(),
2842 actor_id: ToolActorId::new("attacker").unwrap(),
2843 local_id: CapabilityRequestId(1),
2844 };
2845 assert!(matches!(
2846 engine.resolve_approval(ApprovalResolution {
2847 request_key: forged_key,
2848 accepted_sequence: before.requests[0].accepted_sequence,
2849 instance_id: instance(),
2850 generation: generation(1),
2851 decision: ApprovalDecision::ApproveOnce,
2852 }),
2853 Err(ToolEngineError::UnknownRequest { .. })
2854 ));
2855 assert_eq!(engine.snapshot(), before);
2856
2857 assert!(matches!(
2858 engine.resolve_approval(ApprovalResolution {
2859 request_key: request_key(1),
2860 accepted_sequence: before.requests[0].accepted_sequence,
2861 instance_id: AgentInstanceId(999),
2862 generation: generation(1),
2863 decision: ApprovalDecision::ApproveOnce,
2864 }),
2865 Err(ToolEngineError::ApprovalScopeMismatch { .. })
2866 ));
2867 assert_eq!(engine.snapshot(), before);
2868 }
2869
2870 #[test]
2871 fn deactivate_remove_and_reactivate_fail_closed() {
2872 let mut engine = configured(Some(GrantMode::Allow));
2873 engine.request(request(1, 1, 100)).unwrap();
2874 engine.drain_effects();
2875 engine
2876 .set_instance_state(instance(), generation(1), ToolInstanceState::Inactive)
2877 .unwrap();
2878 assert!(engine.snapshot().grants.is_empty());
2879 assert!(matches!(
2880 engine.snapshot().requests[0].status,
2881 CapabilityRequestStatus::InstanceClosed { .. }
2882 ));
2883 assert!(matches!(
2884 engine.drain_completions().completions[0].outcome,
2885 CapabilityTerminalOutcome::InstanceClosed { .. }
2886 ));
2887 assert!(matches!(
2888 engine.resolve_approval(ApprovalResolution {
2889 request_key: request_key(1),
2890 accepted_sequence: engine
2891 .request_snapshot(&request_key(1))
2892 .unwrap()
2893 .accepted_sequence,
2894 instance_id: instance(),
2895 generation: generation(1),
2896 decision: ApprovalDecision::ApproveOnce,
2897 }),
2898 Err(ToolEngineError::InactivePolicyInstance { .. })
2899 ));
2900 assert_eq!(
2901 engine.request(request(2, 1, 100)).unwrap(),
2902 PolicyDecision::Deny(PolicyDenial::InactiveInstance)
2903 );
2904
2905 engine
2906 .set_instance_state(instance(), generation(1), ToolInstanceState::Active)
2907 .unwrap();
2908 engine.set_grant(grant(GrantMode::Allow)).unwrap();
2909 assert_eq!(
2910 engine.request(request(3, 1, 100)).unwrap(),
2911 PolicyDecision::Allow
2912 );
2913 assert!(engine.remove_instance(instance()).unwrap());
2914 assert!(!engine
2915 .snapshot()
2916 .generations
2917 .iter()
2918 .any(|(id, _)| *id == instance()));
2919 assert!(matches!(
2920 engine.resolve_approval(ApprovalResolution {
2921 request_key: request_key(3),
2922 accepted_sequence: engine
2923 .request_snapshot(&request_key(3))
2924 .unwrap()
2925 .accepted_sequence,
2926 instance_id: instance(),
2927 generation: generation(1),
2928 decision: ApprovalDecision::ApproveOnce,
2929 }),
2930 Err(ToolEngineError::UnknownPolicyInstance { .. })
2931 ));
2932 engine.set_generation(instance(), generation(1)).unwrap();
2933 assert!(engine
2934 .snapshot()
2935 .instance_states
2936 .contains(&(instance(), ToolInstanceState::Active)));
2937 }
2938
2939 #[test]
2940 fn explicit_client_close_isolated_to_exact_client() {
2941 let actor_b = ToolActorId::new("consumer.second").unwrap();
2942 let mut engine = configured(Some(GrantMode::RequireApproval));
2943 engine
2944 .set_grant(grant_for_client(
2945 consumer(),
2946 actor_b.clone(),
2947 GrantMode::RequireApproval,
2948 ))
2949 .unwrap();
2950 engine.request(request(1, 1, 100)).unwrap();
2951 engine
2952 .request(request_for_client(1, consumer(), actor_b.clone()))
2953 .unwrap();
2954 assert!(matches!(
2955 engine.close_client(&consumer(), &actor()),
2956 ToolAuthorityOutcome::ClientClosed {
2957 purged_grant_count: 1,
2958 closed_request_count: 1,
2959 }
2960 ));
2961 let snapshot = engine.snapshot();
2962 assert!(snapshot
2963 .grants
2964 .iter()
2965 .any(|grant| grant.key.actor_id == actor_b));
2966 assert!(snapshot.requests.iter().any(|request| {
2967 request.key.actor_id == actor_b
2968 && matches!(request.status, CapabilityRequestStatus::AwaitingApproval)
2969 }));
2970 let completions = engine.drain_completions().completions;
2971 assert_eq!(completions.len(), 1);
2972 assert_eq!(completions[0].request_key.actor_id, actor());
2973 }
2974
2975 #[test]
2976 fn per_client_quota_does_not_block_another_client() {
2977 let actor_b = ToolActorId::new("consumer.second").unwrap();
2978 let mut engine = configured(Some(GrantMode::RequireApproval));
2979 for id in 1..=TOOL_ACTIVE_REQUESTS_PER_CLIENT_MAX as u64 {
2980 engine.request(request(id, 1, 100)).unwrap();
2981 }
2982 assert!(matches!(
2983 engine.request(request(10_000, 1, 100)),
2984 Err(ToolEngineError::ClientRequestCapacityExceeded { .. })
2985 ));
2986 engine
2987 .set_grant(grant_for_client(
2988 consumer(),
2989 actor_b.clone(),
2990 GrantMode::RequireApproval,
2991 ))
2992 .unwrap();
2993 assert_eq!(
2994 engine
2995 .request(request_for_client(1, consumer(), actor_b))
2996 .unwrap(),
2997 PolicyDecision::RequireApproval
2998 );
2999 }
3000
3001 #[test]
3002 fn zero_operation_observation_is_rejected_before_mutation() {
3003 let mut engine = configured(Some(GrantMode::Allow));
3004 engine.request(request(1, 1, 100)).unwrap();
3005 let effect = engine.drain_effects().pop().unwrap();
3006 let before_observation = engine.snapshot();
3007
3008 let mut zero_operation = FakeProvider::succeed(&effect);
3009 zero_operation.operation_id = ToolOperationId(0);
3010 assert!(matches!(
3011 engine.apply_observation(zero_operation),
3012 Err(ToolEngineError::Validation(
3013 ToolValidationError::ZeroIdentifier {
3014 field: "tool operation id"
3015 }
3016 ))
3017 ));
3018 assert_eq!(engine.snapshot(), before_observation);
3019 }
3020
3021 #[test]
3022 fn mutating_expired_approval_consumes_authority_sequence() {
3023 let mut engine = configured(Some(GrantMode::RequireApproval));
3024 engine.request(request(1, 1, 10)).unwrap();
3025 let request_key = request_key(1);
3026 let accepted_sequence = engine
3027 .request_snapshot(&request_key)
3028 .unwrap()
3029 .accepted_sequence;
3030 engine.current_tick = 10;
3031
3032 let outcome = engine
3033 .apply_authority(ToolAuthorityEnvelope {
3034 sequence: 1,
3035 command: ToolAuthorityCommand::ResolveApproval {
3036 resolution: ApprovalResolution {
3037 request_key: request_key.clone(),
3038 accepted_sequence,
3039 instance_id: instance(),
3040 generation: generation(1),
3041 decision: ApprovalDecision::ApproveOnce,
3042 },
3043 },
3044 })
3045 .unwrap();
3046 assert_eq!(
3047 outcome,
3048 ToolAuthorityOutcome::ApprovalExpired {
3049 request_key: request_key.clone(),
3050 accepted_sequence,
3051 }
3052 );
3053 assert!(matches!(
3054 engine.apply_authority(ToolAuthorityEnvelope {
3055 sequence: 1,
3056 command: ToolAuthorityCommand::RevokeGrant {
3057 key: grant(GrantMode::RequireApproval).key,
3058 },
3059 }),
3060 Err(ToolEngineError::AuthoritySequenceRegressed {
3061 current: 1,
3062 requested: 1,
3063 })
3064 ));
3065 assert!(matches!(
3066 engine.request_snapshot(&request_key).unwrap().status,
3067 CapabilityRequestStatus::TimedOut { .. }
3068 ));
3069 assert!(matches!(
3070 engine.drain_completions().completions[0].outcome,
3071 CapabilityTerminalOutcome::TimedOut { .. }
3072 ));
3073 }
3074
3075 #[test]
3076 fn completion_overflow_is_explicit_and_never_blocks_terminal_transition() {
3077 let mut engine = configured(None);
3078 for id in 1..=TOOL_COMPLETIONS_MAX as u64 + 1 {
3079 assert!(matches!(
3080 engine.request(request(id, 1, 100)).unwrap(),
3081 PolicyDecision::Deny(_)
3082 ));
3083 }
3084 assert!(engine
3085 .snapshot()
3086 .requests
3087 .iter()
3088 .all(|request| request.status.is_terminal()));
3089 let batch = engine.drain_completions();
3090 assert_eq!(batch.completions.len(), TOOL_COMPLETIONS_MAX);
3091 assert_eq!(batch.dropped_since_last_drain, 1);
3092 assert_eq!(batch.total_dropped, 1);
3093 let previous_sequence = batch.completions.last().unwrap().sequence;
3094 engine.request(request(50_000, 1, 100)).unwrap();
3095 let next = engine.drain_completions();
3096 assert_eq!(next.completions[0].sequence, previous_sequence + 2);
3097 }
3098
3099 #[test]
3100 fn terminal_key_reuse_is_disambiguated_by_accepted_sequence() {
3101 let mut engine = configured(None);
3102 engine.request(request(1, 1, 100)).unwrap();
3103 let first = engine.drain_completions().completions.pop().unwrap();
3104 engine.request(request(1, 1, 100)).unwrap();
3105 let second = engine.drain_completions().completions.pop().unwrap();
3106 assert_eq!(first.request_key, second.request_key);
3107 assert_ne!(first.accepted_sequence, second.accepted_sequence);
3108 assert!(first.sequence < second.sequence);
3109 assert_eq!(
3110 engine
3111 .request_snapshot(&second.request_key)
3112 .unwrap()
3113 .accepted_sequence,
3114 second.accepted_sequence
3115 );
3116 }
3117
3118 #[test]
3119 fn externally_visible_sequences_never_repeat_after_exhaustion() {
3120 let mut effects = configured(Some(GrantMode::Allow));
3121 effects.next_effect_sequence = u64::MAX;
3122 effects.request(request(1, 1, 100)).unwrap();
3123 let last_effect = effects.drain_effects().pop().unwrap();
3124 assert_eq!(last_effect.sequence, u64::MAX);
3125 assert!(effects.snapshot().effect_sequence_exhausted);
3126 assert!(matches!(
3127 effects.request(request(2, 1, 100)),
3128 Err(ToolEngineError::EffectSequenceExhausted)
3129 ));
3130
3131 let mut completions = configured(None);
3132 completions.next_completion_sequence = u64::MAX;
3133 completions.request(request(1, 1, 100)).unwrap();
3134 let last = completions.drain_completions();
3135 assert_eq!(last.completions[0].sequence, u64::MAX);
3136 assert!(last.sequence_exhausted);
3137 completions.request(request(2, 1, 100)).unwrap();
3138 let dropped = completions.drain_completions();
3139 assert!(dropped.completions.is_empty());
3140 assert_eq!(dropped.dropped_since_last_drain, 1);
3141 assert!(completions
3142 .snapshot()
3143 .requests
3144 .iter()
3145 .all(|request| request.status.is_terminal()));
3146
3147 let mut audit = configured(Some(GrantMode::Allow));
3148 let dropped_before = audit.snapshot().dropped_audit_events;
3149 audit.next_audit_sequence = u64::MAX;
3150 audit.revoke_grant(&grant(GrantMode::Allow).key).unwrap();
3151 audit.set_grant(grant(GrantMode::Allow)).unwrap();
3152 let snapshot = audit.snapshot();
3153 assert_eq!(
3154 snapshot
3155 .audit_events
3156 .iter()
3157 .filter(|event| event.sequence == u64::MAX)
3158 .count(),
3159 1
3160 );
3161 assert!(snapshot.audit_sequence_exhausted);
3162 assert!(snapshot.dropped_audit_events > dropped_before);
3163
3164 audit.revision = u64::MAX;
3165 audit.revision_overflow_count = 0;
3166 audit.advance_time(1).unwrap();
3167 let first = audit.snapshot();
3168 audit.advance_time(2).unwrap();
3169 let second = audit.snapshot();
3170 assert_eq!(first.revision, u64::MAX);
3171 assert_eq!(second.revision, u64::MAX);
3172 assert_eq!(first.revision_overflow_count, 1);
3173 assert_eq!(second.revision_overflow_count, 2);
3174 }
3175}