Skip to main content

gate4agent_tool_engine/
engine.rs

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/// Pure single-writer authority for Gate-owned capability requests.
146///
147/// Callers advance a logical clock and session generations explicitly. The
148/// engine performs no I/O; only drained effects may be handed to a provider.
149#[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    /// Applies the only public policy/approval/client mutation lane.
539    /// Successfully reduced envelopes consume their monotonic authority sequence.
540    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    /// Low-level trusted reducer API. Canonical production ingress must bind
800    /// observations through the kernel/runtime handle before calling it.
801    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    /// Releases bounded terminal outcomes to the requesting shell. Raw inline
1017    /// bytes and opaque provider references exist only in this queue, never in
1018    /// snapshots or audit state.
1019    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(&current_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}