Skip to main content

gate4agent_kernel/
lib.rs

1//! Synchronous host kernel for gate4agent engines.
2
3use gate4agent_catalog::{
4    builtin_registry, resolve_capability_probe_for, resolve_one_shot_plan,
5    resolve_session_option_launch_for, AgentRegistry,
6};
7use gate4agent_engine::Gate4AgentEngine;
8use gate4agent_tool_engine::{ToolEngine, ToolEngineError};
9use gate4agent_tool_protocol::{
10    CapabilityCompletionBatch, CapabilityObservationDisposition, CapabilityProviderDescriptor,
11    CapabilityRequestKey, ObservationIgnoredReason, PolicyDecision, ProviderBindingId,
12    ProviderBoundCapabilityEffectEnvelope, ProviderBoundCapabilityRequest,
13    ProviderRuntimeBindingSnapshot, ProviderRuntimeCommand, ProviderRuntimeEnvelope,
14    ProviderRuntimeSnapshot, ToolAuthorityEnvelope, ToolAuthorityOutcome, ToolEngineSnapshot,
15    ToolInstanceState, ToolProviderId, ToolValidationError,
16};
17use gate4agent_types::{
18    AdapterBinding, AdapterFamily, AgentId, AgentInstanceId, CommandEnvelope, CommandId,
19    ControlCommand, ControlError, ControlEvent, ControlHealth, ControlSnapshot, EffectEnvelope,
20    InputAction, ObservationEnvelope, PipeProtocol, ProviderSource, SessionStatus, TransportKind,
21};
22use std::collections::{BTreeMap, BTreeSet};
23use std::fmt;
24use std::sync::Arc;
25use thiserror::Error;
26
27/// Trusted in-process ingress for the single-writer backend reducer.
28///
29/// Network and IPC shells must authenticate and bind connection-owned
30/// identities before constructing these values. This enum is not itself a
31/// transport authorization boundary.
32#[derive(Clone, Debug, Eq, PartialEq)]
33pub enum BackendIngress {
34    Control(CommandEnvelope),
35    ToolRequest(ProviderBoundCapabilityRequest),
36    ToolAuthority(ToolAuthorityEnvelope),
37    ToolProvider(ProviderRuntimeEnvelope),
38}
39
40#[derive(Clone, Debug, Eq, PartialEq)]
41pub struct ToolRequestOutcome {
42    pub request_key: CapabilityRequestKey,
43    pub accepted_sequence: Option<u64>,
44    pub result: Result<PolicyDecision, KernelToolError>,
45}
46
47#[derive(Clone, Debug, Eq, PartialEq)]
48pub struct ToolAuthorityCommandOutcome {
49    pub sequence: u64,
50    pub result: Result<ToolAuthorityOutcome, KernelToolError>,
51}
52
53#[derive(Clone, Debug, Eq, PartialEq)]
54pub enum BackendIngressOutcome {
55    Control(CommandOutcome),
56    ToolRequest(ToolRequestOutcome),
57    ToolAuthority(ToolAuthorityCommandOutcome),
58    ToolProvider(ProviderRuntimeCommandOutcome),
59}
60
61#[derive(Clone, Debug, Eq, PartialEq)]
62pub struct ProviderRuntimeCommandOutcome {
63    pub sequence: u64,
64    pub binding_id: ProviderBindingId,
65    pub provider_id: ToolProviderId,
66    pub result: Result<ProviderRuntimeTransition, KernelProviderError>,
67}
68
69#[derive(Clone, Debug, Eq, PartialEq)]
70pub enum ProviderRuntimeTransition {
71    Attached,
72    Detached {
73        closed_request_count: usize,
74    },
75    ObservationApplied {
76        operation_id: gate4agent_tool_protocol::ToolOperationId,
77        request_key: CapabilityRequestKey,
78    },
79    ObservationIgnored {
80        operation_id: gate4agent_tool_protocol::ToolOperationId,
81        request_key: CapabilityRequestKey,
82        reason: ObservationIgnoredReason,
83    },
84}
85
86#[derive(Clone, Debug, Eq, Error, PartialEq)]
87pub enum KernelToolError {
88    #[error(transparent)]
89    Engine(#[from] ToolEngineError),
90    #[error(transparent)]
91    Validation(#[from] ToolValidationError),
92    #[error("tool provider '{provider_id}' has no active runtime binding")]
93    ProviderUnavailable { provider_id: ToolProviderId },
94    #[error(
95        "tool request targets provider '{provider_id}' binding {requested:?}, current binding is {current:?}"
96    )]
97    ProviderBindingMismatch {
98        provider_id: ToolProviderId,
99        current: ProviderBindingId,
100        requested: Option<ProviderBindingId>,
101    },
102    #[error("tool lane is blocked by a control/tool integration failure")]
103    IntegrationBlocked,
104}
105
106#[derive(Clone, Debug, Eq, Error, PartialEq)]
107pub enum KernelProviderError {
108    #[error(transparent)]
109    Validation(#[from] gate4agent_tool_protocol::ToolValidationError),
110    #[error("tool provider runtime sequence is exhausted")]
111    SequenceExhausted,
112    #[error("tool provider runtime sequence regressed from {current} to {requested}")]
113    SequenceRegressed { current: u64, requested: u64 },
114    #[error("tool provider '{provider_id}' is not registered")]
115    UnknownProvider { provider_id: ToolProviderId },
116    #[error("attach binding {binding_id:?} must equal provider runtime sequence {sequence}")]
117    InvalidAttachBinding {
118        sequence: u64,
119        binding_id: ProviderBindingId,
120    },
121    #[error("tool provider '{provider_id}' is already attached as binding {binding_id:?}")]
122    AlreadyAttached {
123        provider_id: ToolProviderId,
124        binding_id: ProviderBindingId,
125    },
126    #[error("tool provider '{provider_id}' has no active runtime binding")]
127    NotAttached { provider_id: ToolProviderId },
128    #[error("tool provider '{provider_id}' is attached as {current:?}, not {requested:?}")]
129    BindingMismatch {
130        provider_id: ToolProviderId,
131        current: ProviderBindingId,
132        requested: ProviderBindingId,
133    },
134    #[error(transparent)]
135    Engine(#[from] ToolEngineError),
136    #[error("tool provider lane is blocked by a control/tool integration failure")]
137    IntegrationBlocked,
138}
139
140#[derive(Clone, Debug, Eq, Error, PartialEq)]
141pub enum KernelIntegrationError {
142    #[error("kernel logical tick exhausted at {current_tick}")]
143    LogicalTickExhausted { current_tick: u64 },
144    #[error("kernel snapshot revision exhausted at {current_revision}")]
145    BackendRevisionExhausted { current_revision: u64 },
146    #[error("control engine entered terminal counter exhaustion: {health:?}")]
147    ControlHealthExhausted { health: ControlHealth },
148    #[error("control observation for {instance_id:?} failed: {source}")]
149    ControlObservation {
150        instance_id: AgentInstanceId,
151        #[source]
152        source: ControlError,
153    },
154    #[error("tool clock advance failed: {source}")]
155    ToolClock {
156        #[source]
157        source: ToolEngineError,
158    },
159    #[error("tool instance sync failed for {instance_id:?}: {source}")]
160    ToolInstanceSync {
161        instance_id: AgentInstanceId,
162        #[source]
163        source: ToolEngineError,
164    },
165    #[error(
166        "tool effect {operation_id:?} targets provider '{provider_id}' without an active runtime binding"
167    )]
168    ToolEffectProviderUnbound {
169        operation_id: gate4agent_tool_protocol::ToolOperationId,
170        provider_id: ToolProviderId,
171    },
172}
173
174#[derive(Clone, Debug, Eq, PartialEq)]
175pub struct BackendSnapshot {
176    pub revision: u64,
177    pub logical_tick: u64,
178    pub control: Arc<ControlSnapshot>,
179    pub tools: Arc<ToolEngineSnapshot>,
180    pub provider_runtime: ProviderRuntimeSnapshot,
181}
182
183impl Default for BackendSnapshot {
184    fn default() -> Self {
185        Self {
186            revision: 0,
187            logical_tick: 0,
188            control: Arc::new(ControlSnapshot {
189                revision: 0,
190                health: ControlHealth::default(),
191                sessions: Vec::new(),
192            }),
193            tools: Arc::new(ToolEngine::new().snapshot()),
194            provider_runtime: ProviderRuntimeSnapshot {
195                last_sequence: 0,
196                sequence_exhausted: false,
197                bindings: Vec::new(),
198            },
199        }
200    }
201}
202
203#[derive(Clone, Debug, Eq, PartialEq)]
204pub struct CommandOutcome {
205    pub command_id: CommandId,
206    pub result: Result<(), KernelCommandError>,
207}
208
209#[derive(Clone, Debug, Eq, Error, PartialEq)]
210pub enum KernelCommandError {
211    #[error("kernel control plane is blocked: {reason}")]
212    IntegrationBlocked { reason: KernelIntegrationError },
213    #[error("agent '{agent_id}' is not present in the kernel catalog")]
214    UnknownAgent { agent_id: AgentId },
215    #[error("agent '{agent_id}' does not declare capability '{capability}'")]
216    UnsupportedCapability {
217        agent_id: AgentId,
218        capability: &'static str,
219    },
220    #[error("agent '{agent_id}' does not support transport {transport:?}")]
221    UnsupportedTransport {
222        agent_id: AgentId,
223        transport: TransportKind,
224    },
225    #[error(
226        "agent '{agent_id}' does not declare {family:?} provider source '{adapter_id}' at revision '{revision}'"
227    )]
228    InvalidProviderSource {
229        agent_id: AgentId,
230        family: AdapterFamily,
231        adapter_id: String,
232        revision: String,
233    },
234    #[error("agent '{agent_id}' session options are invalid: {message}")]
235    InvalidSessionOptions { agent_id: AgentId, message: String },
236    #[error("agent '{agent_id}' capability probe is invalid: {message}")]
237    InvalidCapabilityProbe { agent_id: AgentId, message: String },
238    #[error(transparent)]
239    Control(#[from] ControlError),
240}
241
242impl KernelCommandError {
243    /// True for exactly the `UnsupportedTransport` variant. Exists so a
244    /// caller several layers up the stack -- `gate4agent-node`'s spawn-
245    /// dispatch waiter, specifically, which cannot and should not depend on
246    /// this crate directly (see this crate's own `Forbidden` list) -- can
247    /// classify one specific, already-well-known kernel rejection by
248    /// calling this inherent method on the value it already has (via
249    /// `gate4agent-runtime-native`'s `CommandOutcome` re-export), without
250    /// naming `KernelCommandError`'s own type path or matching its full
251    /// variant set.
252    pub fn is_unsupported_transport(&self) -> bool {
253        matches!(self, Self::UnsupportedTransport { .. })
254    }
255}
256
257#[derive(Clone, Debug, Eq, PartialEq)]
258pub struct KernelStep {
259    pub command_outcomes: Vec<CommandOutcome>,
260    pub effects: Vec<EffectEnvelope>,
261    pub snapshot: ControlSnapshot,
262    pub events: Vec<ControlEvent>,
263    pub ingress_outcomes: Vec<BackendIngressOutcome>,
264    pub tool_effects: Vec<ProviderBoundCapabilityEffectEnvelope>,
265    pub tool_completions: CapabilityCompletionBatch,
266    pub backend_snapshot: BackendSnapshot,
267    pub integration_errors: Vec<KernelIntegrationError>,
268}
269
270/// Owns the provider catalog and the single session-state writer.
271///
272/// Phase order is fixed: clock advance, ordered ingress, control observations,
273/// tool observations, effect/completion drains, atomic snapshot, then ordered
274/// event drain. No phase awaits or performs external work. The kernel is
275/// intentionally not cloneable: cloning would fork provider binding and
276/// sequence authority while preserving otherwise valid runtime identities.
277pub struct Gate4AgentKernel {
278    catalog: AgentRegistry,
279    engine: Gate4AgentEngine,
280    tool_engine: ToolEngine,
281    provider_bindings: BTreeMap<ToolProviderId, ProviderBindingId>,
282    last_provider_sequence: u64,
283    provider_sequence_exhausted: bool,
284    logical_tick: u64,
285    backend_revision: u64,
286}
287
288impl fmt::Debug for Gate4AgentKernel {
289    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
290        formatter
291            .debug_struct("Gate4AgentKernel")
292            .field("catalog", &self.catalog)
293            .field("engine", &self.engine)
294            .field("tools", &self.tool_engine.snapshot())
295            .field("provider_runtime", &self.provider_runtime_snapshot())
296            .field("logical_tick", &self.logical_tick)
297            .field("backend_revision", &self.backend_revision)
298            .finish()
299    }
300}
301
302impl Gate4AgentKernel {
303    pub fn new(catalog: AgentRegistry) -> Self {
304        Self {
305            catalog,
306            engine: Gate4AgentEngine::new(),
307            tool_engine: ToolEngine::new(),
308            provider_bindings: BTreeMap::new(),
309            last_provider_sequence: 0,
310            provider_sequence_exhausted: false,
311            logical_tick: 0,
312            backend_revision: 0,
313        }
314    }
315
316    pub fn with_tool_providers(
317        catalog: AgentRegistry,
318        providers: impl IntoIterator<Item = CapabilityProviderDescriptor>,
319    ) -> Result<Self, ToolEngineError> {
320        let mut kernel = Self::new(catalog);
321        for provider in providers {
322            kernel.tool_engine.register_provider(provider)?;
323        }
324        Ok(kernel)
325    }
326
327    pub fn with_builtin_catalog() -> Self {
328        Self::new(builtin_registry().clone())
329    }
330
331    pub fn step(
332        &mut self,
333        commands: impl IntoIterator<Item = CommandEnvelope>,
334        observations: impl IntoIterator<Item = ObservationEnvelope>,
335    ) -> KernelStep {
336        self.step_control_plane(
337            commands.into_iter().map(BackendIngress::Control),
338            observations,
339        )
340    }
341
342    /// Reduces one backend tick without performing external work.
343    ///
344    /// The phase order is fixed: advance the kernel-owned logical clock;
345    /// reduce ordered control/tool ingress; reduce control observations while
346    /// synchronizing the affected tool instance after every observation;
347    /// reduce provider-bound tool observations; then drain effects,
348    /// completions, snapshots, and events. A tool integration failure blocks
349    /// tool ingress and effect release for the rest of the tick.
350    pub fn step_control_plane(
351        &mut self,
352        ingress: impl IntoIterator<Item = BackendIngress>,
353        control_observations: impl IntoIterator<Item = ObservationEnvelope>,
354    ) -> KernelStep {
355        let ingress = ingress.into_iter().collect::<Vec<_>>();
356        let control_observations = control_observations.into_iter().collect::<Vec<_>>();
357
358        if let Err(error) = self.advance_backend_clock() {
359            return self.blocked_step(ingress, error);
360        }
361
362        let mut command_outcomes = Vec::new();
363        let mut ingress_outcomes = Vec::new();
364        let mut integration_errors = Vec::new();
365        let mut control_lane_open = true;
366        let mut tool_lane_open = true;
367        self.block_on_control_health(
368            &mut control_lane_open,
369            &mut tool_lane_open,
370            &mut integration_errors,
371        );
372        if control_lane_open {
373            self.reconcile_or_block(&mut tool_lane_open, &mut integration_errors);
374        }
375
376        for item in ingress {
377            match item {
378                BackendIngress::Control(command) => {
379                    let command_id = command.id;
380                    let instance_id = command.command.instance_id();
381                    let attempted = control_lane_open;
382                    let result = if attempted {
383                        self.apply_validated_command(command)
384                    } else {
385                        Err(KernelCommandError::IntegrationBlocked {
386                            reason: self
387                                .control_health_error()
388                                .expect("closed control lane has terminal health"),
389                        })
390                    };
391                    if attempted {
392                        if let Err(error) = &result {
393                            self.engine.record_command_rejection(
394                                command_id,
395                                instance_id,
396                                error.to_string(),
397                            );
398                        }
399                    }
400                    let outcome = CommandOutcome { command_id, result };
401                    command_outcomes.push(outcome.clone());
402                    ingress_outcomes.push(BackendIngressOutcome::Control(outcome));
403                    if attempted {
404                        self.sync_or_block(
405                            instance_id,
406                            &mut tool_lane_open,
407                            &mut integration_errors,
408                        );
409                        self.block_on_control_health(
410                            &mut control_lane_open,
411                            &mut tool_lane_open,
412                            &mut integration_errors,
413                        );
414                    }
415                }
416                BackendIngress::ToolRequest(bound_request) => {
417                    let request_key = bound_request.key();
418                    let provider_id = bound_request.request.request.provider_id.clone();
419                    let result = if !tool_lane_open {
420                        Err(KernelToolError::IntegrationBlocked)
421                    } else if let Err(error) = bound_request.validate_provider_binding() {
422                        Err(KernelToolError::Validation(error))
423                    } else if self.tool_engine.provider_exists(&provider_id) {
424                        match self.provider_bindings.get(&provider_id).copied() {
425                            None => Err(KernelToolError::ProviderUnavailable { provider_id }),
426                            Some(current) if bound_request.provider_binding_id != Some(current) => {
427                                Err(KernelToolError::ProviderBindingMismatch {
428                                    provider_id,
429                                    current,
430                                    requested: bound_request.provider_binding_id,
431                                })
432                            }
433                            Some(_) => self
434                                .tool_engine
435                                .request(bound_request.request)
436                                .map_err(Into::into),
437                        }
438                    } else {
439                        self.tool_engine
440                            .request(bound_request.request)
441                            .map_err(Into::into)
442                    };
443                    let accepted_sequence = result.as_ref().ok().and_then(|_| {
444                        self.tool_engine
445                            .request_snapshot(&request_key)
446                            .map(|snapshot| snapshot.accepted_sequence)
447                    });
448                    ingress_outcomes.push(BackendIngressOutcome::ToolRequest(ToolRequestOutcome {
449                        request_key,
450                        accepted_sequence,
451                        result,
452                    }));
453                }
454                BackendIngress::ToolAuthority(authority) => {
455                    let sequence = authority.sequence;
456                    let result = if tool_lane_open {
457                        self.tool_engine
458                            .apply_authority(authority)
459                            .map_err(Into::into)
460                    } else {
461                        Err(KernelToolError::IntegrationBlocked)
462                    };
463                    ingress_outcomes.push(BackendIngressOutcome::ToolAuthority(
464                        ToolAuthorityCommandOutcome { sequence, result },
465                    ));
466                }
467                BackendIngress::ToolProvider(envelope) => {
468                    let outcome = self.apply_provider_runtime(envelope, tool_lane_open);
469                    ingress_outcomes.push(BackendIngressOutcome::ToolProvider(outcome));
470                }
471            }
472        }
473
474        for observation in control_observations {
475            let instance_id = observation.instance_id;
476            match self.engine.try_apply_observation(observation) {
477                Ok(()) => {
478                    self.sync_or_block(instance_id, &mut tool_lane_open, &mut integration_errors)
479                }
480                Err(source) => {
481                    control_lane_open = false;
482                    tool_lane_open = false;
483                    integration_errors.push(KernelIntegrationError::ControlObservation {
484                        instance_id,
485                        source,
486                    });
487                }
488            }
489            self.block_on_control_health(
490                &mut control_lane_open,
491                &mut tool_lane_open,
492                &mut integration_errors,
493            );
494        }
495
496        let effects = self.engine.drain_effects();
497        let tool_effects = if tool_lane_open {
498            match self.drain_bound_tool_effects() {
499                Ok(effects) => effects,
500                Err(error) => {
501                    integration_errors.push(error);
502                    Vec::new()
503                }
504            }
505        } else {
506            Vec::new()
507        };
508        let tool_completions = self.tool_engine.drain_completions();
509        let backend_snapshot = self.backend_snapshot();
510        let snapshot = (*backend_snapshot.control).clone();
511        let events = self.engine.drain_events();
512
513        KernelStep {
514            command_outcomes,
515            effects,
516            snapshot,
517            events,
518            ingress_outcomes,
519            tool_effects,
520            tool_completions,
521            backend_snapshot,
522            integration_errors,
523        }
524    }
525
526    pub fn snapshot(&self) -> ControlSnapshot {
527        self.engine.snapshot()
528    }
529
530    pub fn tool_snapshot(&self) -> ToolEngineSnapshot {
531        self.tool_engine.snapshot()
532    }
533
534    pub fn backend_snapshot(&self) -> BackendSnapshot {
535        BackendSnapshot {
536            revision: self.backend_revision,
537            logical_tick: self.logical_tick,
538            control: Arc::new(self.engine.snapshot()),
539            tools: Arc::new(self.tool_engine.snapshot()),
540            provider_runtime: self.provider_runtime_snapshot(),
541        }
542    }
543
544    pub fn catalog(&self) -> &AgentRegistry {
545        &self.catalog
546    }
547
548    fn provider_runtime_snapshot(&self) -> ProviderRuntimeSnapshot {
549        ProviderRuntimeSnapshot {
550            last_sequence: self.last_provider_sequence,
551            sequence_exhausted: self.provider_sequence_exhausted,
552            bindings: self
553                .provider_bindings
554                .iter()
555                .map(|(provider_id, binding_id)| ProviderRuntimeBindingSnapshot {
556                    binding_id: *binding_id,
557                    provider_id: provider_id.clone(),
558                })
559                .collect(),
560        }
561    }
562
563    fn apply_provider_runtime(
564        &mut self,
565        envelope: ProviderRuntimeEnvelope,
566        tool_lane_open: bool,
567    ) -> ProviderRuntimeCommandOutcome {
568        let sequence = envelope.sequence;
569        let (binding_id, provider_id) = provider_runtime_subject(&envelope.command);
570        let mut result = self.reduce_provider_runtime(envelope, tool_lane_open);
571        if self.provider_sequence_exhausted {
572            if let Err(error) = self.retire_exhausted_provider_bindings() {
573                result = Err(KernelProviderError::Engine(error));
574            }
575        }
576        ProviderRuntimeCommandOutcome {
577            sequence,
578            binding_id,
579            provider_id,
580            result,
581        }
582    }
583
584    fn retire_exhausted_provider_bindings(&mut self) -> Result<(), ToolEngineError> {
585        let provider_ids = self.provider_bindings.keys().cloned().collect::<Vec<_>>();
586        for provider_id in provider_ids {
587            self.tool_engine.detach_provider_runtime(&provider_id)?;
588            self.provider_bindings.remove(&provider_id);
589        }
590        Ok(())
591    }
592
593    fn reduce_provider_runtime(
594        &mut self,
595        envelope: ProviderRuntimeEnvelope,
596        tool_lane_open: bool,
597    ) -> Result<ProviderRuntimeTransition, KernelProviderError> {
598        envelope.validate()?;
599        if self.provider_sequence_exhausted {
600            return Err(KernelProviderError::SequenceExhausted);
601        }
602        if envelope.sequence <= self.last_provider_sequence {
603            return Err(KernelProviderError::SequenceRegressed {
604                current: self.last_provider_sequence,
605                requested: envelope.sequence,
606            });
607        }
608
609        self.last_provider_sequence = envelope.sequence;
610        if envelope.sequence == u64::MAX {
611            self.provider_sequence_exhausted = true;
612        }
613        if !tool_lane_open {
614            return Err(KernelProviderError::IntegrationBlocked);
615        }
616
617        match envelope.command {
618            ProviderRuntimeCommand::Attach {
619                binding_id,
620                provider_id,
621            } => {
622                if binding_id.0 != envelope.sequence {
623                    return Err(KernelProviderError::InvalidAttachBinding {
624                        sequence: envelope.sequence,
625                        binding_id,
626                    });
627                }
628                if !self.tool_engine.provider_exists(&provider_id) {
629                    return Err(KernelProviderError::UnknownProvider { provider_id });
630                }
631                if let Some(current) = self.provider_bindings.get(&provider_id).copied() {
632                    return Err(KernelProviderError::AlreadyAttached {
633                        provider_id,
634                        binding_id: current,
635                    });
636                }
637                self.provider_bindings.insert(provider_id, binding_id);
638                Ok(ProviderRuntimeTransition::Attached)
639            }
640            ProviderRuntimeCommand::Detach {
641                binding_id,
642                provider_id,
643            } => {
644                if !self.tool_engine.provider_exists(&provider_id) {
645                    return Err(KernelProviderError::UnknownProvider { provider_id });
646                }
647                self.require_provider_binding(&provider_id, binding_id)?;
648                let closed_request_count =
649                    self.tool_engine.detach_provider_runtime(&provider_id)?;
650                self.provider_bindings.remove(&provider_id);
651                Ok(ProviderRuntimeTransition::Detached {
652                    closed_request_count,
653                })
654            }
655            ProviderRuntimeCommand::Observe {
656                binding_id,
657                observation,
658            } => {
659                let provider_id = observation.provider_id.clone();
660                if !self.tool_engine.provider_exists(&provider_id) {
661                    return Err(KernelProviderError::UnknownProvider { provider_id });
662                }
663                self.require_provider_binding(&provider_id, binding_id)?;
664                let operation_id = observation.operation_id;
665                let request_key = observation.request_key.clone();
666                match self.tool_engine.apply_observation(observation)? {
667                    CapabilityObservationDisposition::Applied => {
668                        Ok(ProviderRuntimeTransition::ObservationApplied {
669                            operation_id,
670                            request_key,
671                        })
672                    }
673                    CapabilityObservationDisposition::Ignored { reason } => {
674                        Ok(ProviderRuntimeTransition::ObservationIgnored {
675                            operation_id,
676                            request_key,
677                            reason,
678                        })
679                    }
680                }
681            }
682        }
683    }
684
685    fn require_provider_binding(
686        &self,
687        provider_id: &ToolProviderId,
688        requested: ProviderBindingId,
689    ) -> Result<(), KernelProviderError> {
690        let Some(current) = self.provider_bindings.get(provider_id).copied() else {
691            return Err(KernelProviderError::NotAttached {
692                provider_id: provider_id.clone(),
693            });
694        };
695        if current != requested {
696            return Err(KernelProviderError::BindingMismatch {
697                provider_id: provider_id.clone(),
698                current,
699                requested,
700            });
701        }
702        Ok(())
703    }
704
705    fn drain_bound_tool_effects(
706        &mut self,
707    ) -> Result<Vec<ProviderBoundCapabilityEffectEnvelope>, KernelIntegrationError> {
708        let effects = self.tool_engine.drain_effects();
709        let mut bound = Vec::with_capacity(effects.len());
710        for effect in effects {
711            let Some(binding_id) = self.provider_bindings.get(&effect.provider_id).copied() else {
712                return Err(KernelIntegrationError::ToolEffectProviderUnbound {
713                    operation_id: effect.operation_id,
714                    provider_id: effect.provider_id,
715                });
716            };
717            bound.push(ProviderBoundCapabilityEffectEnvelope { binding_id, effect });
718        }
719        Ok(bound)
720    }
721
722    fn advance_backend_clock(&mut self) -> Result<(), KernelIntegrationError> {
723        let next_tick = self.logical_tick.checked_add(1).ok_or(
724            KernelIntegrationError::LogicalTickExhausted {
725                current_tick: self.logical_tick,
726            },
727        )?;
728        let next_revision = self.backend_revision.checked_add(1).ok_or(
729            KernelIntegrationError::BackendRevisionExhausted {
730                current_revision: self.backend_revision,
731            },
732        )?;
733        self.tool_engine
734            .advance_time(next_tick)
735            .map_err(|source| KernelIntegrationError::ToolClock { source })?;
736        self.logical_tick = next_tick;
737        self.backend_revision = next_revision;
738        Ok(())
739    }
740
741    fn sync_or_block(
742        &mut self,
743        instance_id: AgentInstanceId,
744        tool_lane_open: &mut bool,
745        integration_errors: &mut Vec<KernelIntegrationError>,
746    ) {
747        if let Err(error) = self.sync_control_instance(instance_id) {
748            *tool_lane_open = false;
749            integration_errors.push(error);
750        }
751    }
752
753    fn control_health_error(&self) -> Option<KernelIntegrationError> {
754        terminal_control_health_error(self.engine.health())
755    }
756
757    fn block_on_control_health(
758        &self,
759        control_lane_open: &mut bool,
760        tool_lane_open: &mut bool,
761        integration_errors: &mut Vec<KernelIntegrationError>,
762    ) {
763        block_lanes_on_control_health(
764            self.engine.health(),
765            control_lane_open,
766            tool_lane_open,
767            integration_errors,
768        );
769    }
770
771    fn reconcile_or_block(
772        &mut self,
773        tool_lane_open: &mut bool,
774        integration_errors: &mut Vec<KernelIntegrationError>,
775    ) {
776        let control_instance_ids = self.engine.session_instance_ids().collect::<BTreeSet<_>>();
777        for instance_id in &control_instance_ids {
778            self.sync_or_block(*instance_id, tool_lane_open, integration_errors);
779        }
780        let tool_instance_ids = self.tool_engine.instance_ids().collect::<Vec<_>>();
781        for instance_id in tool_instance_ids {
782            if !control_instance_ids.contains(&instance_id) {
783                self.sync_or_block(instance_id, tool_lane_open, integration_errors);
784            }
785        }
786    }
787
788    fn sync_control_instance(
789        &mut self,
790        instance_id: AgentInstanceId,
791    ) -> Result<(), KernelIntegrationError> {
792        let Some(session) = self.engine.session_snapshot(instance_id).cloned() else {
793            self.tool_engine
794                .remove_instance(instance_id)
795                .map_err(|source| KernelIntegrationError::ToolInstanceSync {
796                    instance_id,
797                    source,
798                })?;
799            return Ok(());
800        };
801
802        self.tool_engine
803            .set_generation(instance_id, session.generation)
804            .map_err(|source| KernelIntegrationError::ToolInstanceSync {
805                instance_id,
806                source,
807            })?;
808        let state = if session.status == SessionStatus::Running {
809            ToolInstanceState::Active
810        } else {
811            ToolInstanceState::Inactive
812        };
813        self.tool_engine
814            .set_instance_state(instance_id, session.generation, state)
815            .map_err(|source| KernelIntegrationError::ToolInstanceSync {
816                instance_id,
817                source,
818            })
819    }
820
821    fn blocked_step(
822        &self,
823        ingress: Vec<BackendIngress>,
824        reason: KernelIntegrationError,
825    ) -> KernelStep {
826        let mut command_outcomes = Vec::new();
827        let mut ingress_outcomes = Vec::new();
828        for item in ingress {
829            match item {
830                BackendIngress::Control(command) => {
831                    let outcome = CommandOutcome {
832                        command_id: command.id,
833                        result: Err(KernelCommandError::IntegrationBlocked {
834                            reason: reason.clone(),
835                        }),
836                    };
837                    command_outcomes.push(outcome.clone());
838                    ingress_outcomes.push(BackendIngressOutcome::Control(outcome));
839                }
840                BackendIngress::ToolRequest(request) => {
841                    ingress_outcomes.push(BackendIngressOutcome::ToolRequest(ToolRequestOutcome {
842                        request_key: request.key(),
843                        accepted_sequence: None,
844                        result: Err(KernelToolError::IntegrationBlocked),
845                    }));
846                }
847                BackendIngress::ToolAuthority(authority) => {
848                    ingress_outcomes.push(BackendIngressOutcome::ToolAuthority(
849                        ToolAuthorityCommandOutcome {
850                            sequence: authority.sequence,
851                            result: Err(KernelToolError::IntegrationBlocked),
852                        },
853                    ));
854                }
855                BackendIngress::ToolProvider(envelope) => {
856                    let (binding_id, provider_id) = provider_runtime_subject(&envelope.command);
857                    ingress_outcomes.push(BackendIngressOutcome::ToolProvider(
858                        ProviderRuntimeCommandOutcome {
859                            sequence: envelope.sequence,
860                            binding_id,
861                            provider_id,
862                            result: Err(KernelProviderError::IntegrationBlocked),
863                        },
864                    ));
865                }
866            }
867        }
868        let backend_snapshot = self.backend_snapshot();
869        let snapshot = (*backend_snapshot.control).clone();
870        let tool_completions = CapabilityCompletionBatch {
871            completions: Vec::new(),
872            dropped_since_last_drain: 0,
873            total_dropped: backend_snapshot.tools.dropped_completions,
874            next_sequence: backend_snapshot.tools.next_completion_sequence,
875            sequence_exhausted: backend_snapshot.tools.completion_sequence_exhausted,
876        };
877
878        KernelStep {
879            command_outcomes,
880            effects: Vec::new(),
881            snapshot,
882            events: Vec::new(),
883            ingress_outcomes,
884            tool_effects: Vec::new(),
885            tool_completions,
886            backend_snapshot,
887            integration_errors: vec![reason],
888        }
889    }
890
891    fn apply_validated_command(
892        &mut self,
893        mut command: CommandEnvelope,
894    ) -> Result<(), KernelCommandError> {
895        if let ControlCommand::Register {
896            agent_id,
897            transport,
898            ..
899        } = &command.command
900        {
901            let Some(spec) = self.catalog.get(agent_id) else {
902                return Err(KernelCommandError::UnknownAgent {
903                    agent_id: agent_id.clone(),
904                });
905            };
906            let supported = match transport {
907                TransportKind::Pty => spec.capabilities.transports.pty,
908                TransportKind::Pipe => spec.capabilities.transports.pipe.is_some(),
909                TransportKind::Acp => spec.capabilities.transports.acp.is_some(),
910            };
911            if !supported {
912                return Err(KernelCommandError::UnsupportedTransport {
913                    agent_id: agent_id.clone(),
914                    transport: *transport,
915                });
916            }
917        }
918        if let ControlCommand::SendInput {
919            instance_id,
920            action: InputAction::AgentCommand(_),
921        } = &command.command
922        {
923            if let Some(session) = self.engine.session_snapshot(*instance_id) {
924                let supports_agent_commands = self
925                    .catalog
926                    .get(&session.agent_id)
927                    .is_some_and(|spec| spec.capabilities.agent_commands.is_some());
928                if !supports_agent_commands {
929                    return Err(KernelCommandError::UnsupportedCapability {
930                        agent_id: session.agent_id.clone(),
931                        capability: "agent-commands",
932                    });
933                }
934            }
935        }
936        if let ControlCommand::DiscoverHistory { instance_id, .. }
937        | ControlCommand::LoadHistory { instance_id, .. } = &command.command
938        {
939            if let Some(session) = self.engine.session_snapshot(*instance_id) {
940                let supports_history = self
941                    .catalog
942                    .get(&session.agent_id)
943                    .is_some_and(|spec| spec.capabilities.adapters.history.is_some());
944                if !supports_history {
945                    return Err(KernelCommandError::UnsupportedCapability {
946                        agent_id: session.agent_id.clone(),
947                        capability: "history",
948                    });
949                }
950            }
951        }
952        if let ControlCommand::ProbeCapabilities { instance_id, .. } = &command.command {
953            if let Some(session) = self.engine.session_snapshot(*instance_id) {
954                let spec = self
955                    .catalog
956                    .get(&session.agent_id)
957                    .expect("registered agent must remain in kernel catalog");
958                if spec.capabilities.adapters.capability_probe.is_none() {
959                    return Err(KernelCommandError::UnsupportedCapability {
960                        agent_id: session.agent_id.clone(),
961                        capability: "capability-probe",
962                    });
963                }
964                resolve_capability_probe_for(spec).map_err(|error| {
965                    KernelCommandError::InvalidCapabilityProbe {
966                        agent_id: session.agent_id.clone(),
967                        message: error.to_string(),
968                    }
969                })?;
970            }
971        }
972        if let ControlCommand::Resume { instance_id, .. } = &command.command {
973            if let Some(session) = self.engine.session_snapshot(*instance_id) {
974                let supports_resume = self
975                    .catalog
976                    .get(&session.agent_id)
977                    .is_some_and(|spec| spec.capabilities.adapters.resume.is_some());
978                if !supports_resume
979                    || !matches!(session.transport, TransportKind::Pty | TransportKind::Pipe)
980                {
981                    return Err(KernelCommandError::UnsupportedCapability {
982                        agent_id: session.agent_id.clone(),
983                        capability: "resume",
984                    });
985                }
986            }
987        }
988        if let ControlCommand::Start {
989            instance_id,
990            request,
991            ..
992        } = &mut command.command
993        {
994            if let Some(session) = self.engine.session_snapshot(*instance_id) {
995                let spec = self
996                    .catalog
997                    .get(&session.agent_id)
998                    .expect("registered agent must remain in kernel catalog");
999                let one_shot = spec
1000                    .capabilities
1001                    .transports
1002                    .pipe
1003                    .as_ref()
1004                    .filter(|transport| transport.protocol == PipeProtocol::OneShotText);
1005                if session.transport == TransportKind::Pipe && one_shot.is_some() {
1006                    let prompt = request.initial_prompt.as_deref().unwrap_or_default();
1007                    let binding = spec
1008                        .capabilities
1009                        .adapters
1010                        .one_shot
1011                        .as_ref()
1012                        .expect("validated one-shot transport binding");
1013                    let resolved = resolve_one_shot_plan(
1014                        &binding.id,
1015                        &spec.launch,
1016                        prompt,
1017                        request.session_options.as_ref(),
1018                    )
1019                    .map_err(|error| {
1020                        KernelCommandError::InvalidSessionOptions {
1021                            agent_id: session.agent_id.clone(),
1022                            message: error.to_string(),
1023                        }
1024                    })?;
1025                    request.session_options = Some(resolved.applied);
1026                } else if let Some(session_options) = &request.session_options {
1027                    if session.transport != TransportKind::Pty
1028                        || spec.capabilities.adapters.session_options.is_none()
1029                    {
1030                        return Err(KernelCommandError::UnsupportedCapability {
1031                            agent_id: session.agent_id.clone(),
1032                            capability: "pty-session-options",
1033                        });
1034                    }
1035                    let resolved = resolve_session_option_launch_for(spec, session_options, &[])
1036                        .map_err(|error| KernelCommandError::InvalidSessionOptions {
1037                            agent_id: session.agent_id.clone(),
1038                            message: error.to_string(),
1039                        })?;
1040                    request.session_options = resolved.applied;
1041                }
1042            }
1043        }
1044        if let ControlCommand::IngestProvider {
1045            instance_id,
1046            source,
1047            ..
1048        } = &command.command
1049        {
1050            if let Some(session) = self.engine.session_snapshot(*instance_id) {
1051                let spec = self
1052                    .catalog
1053                    .get(&session.agent_id)
1054                    .expect("registered agent must remain in kernel catalog");
1055                if declared_provider_binding(spec, source) != Some(&source.binding) {
1056                    return Err(KernelCommandError::InvalidProviderSource {
1057                        agent_id: session.agent_id.clone(),
1058                        family: source.family,
1059                        adapter_id: source.binding.id.to_string(),
1060                        revision: source.binding.revision.clone(),
1061                    });
1062                }
1063            }
1064        }
1065        self.engine.apply_command(command).map_err(Into::into)
1066    }
1067}
1068
1069fn provider_runtime_subject(
1070    command: &ProviderRuntimeCommand,
1071) -> (ProviderBindingId, ToolProviderId) {
1072    match command {
1073        ProviderRuntimeCommand::Attach {
1074            binding_id,
1075            provider_id,
1076        }
1077        | ProviderRuntimeCommand::Detach {
1078            binding_id,
1079            provider_id,
1080        } => (*binding_id, provider_id.clone()),
1081        ProviderRuntimeCommand::Observe {
1082            binding_id,
1083            observation,
1084        } => (*binding_id, observation.provider_id.clone()),
1085    }
1086}
1087
1088fn terminal_control_health_error(health: ControlHealth) -> Option<KernelIntegrationError> {
1089    (health.operation_id_exhausted
1090        || health.event_sequence_exhausted
1091        || health.revision_exhausted
1092        || health.provider_sequence_exhausted_sessions > 0)
1093        .then_some(KernelIntegrationError::ControlHealthExhausted { health })
1094}
1095
1096fn block_lanes_on_control_health(
1097    health: ControlHealth,
1098    control_lane_open: &mut bool,
1099    tool_lane_open: &mut bool,
1100    integration_errors: &mut Vec<KernelIntegrationError>,
1101) {
1102    let Some(error) = terminal_control_health_error(health) else {
1103        return;
1104    };
1105    *control_lane_open = false;
1106    *tool_lane_open = false;
1107    if !integration_errors.contains(&error) {
1108        integration_errors.push(error);
1109    }
1110}
1111
1112fn declared_provider_binding<'a>(
1113    spec: &'a gate4agent_catalog::AgentSpec,
1114    source: &ProviderSource,
1115) -> Option<&'a AdapterBinding> {
1116    match source.family {
1117        AdapterFamily::PtySemantic => spec.capabilities.transports.pty_adapter.as_ref(),
1118        AdapterFamily::Pipe => {
1119            let transport_binding = spec
1120                .capabilities
1121                .transports
1122                .pipe
1123                .as_ref()
1124                .map(|transport| &transport.adapter);
1125            transport_binding
1126                .filter(|binding| *binding == &source.binding)
1127                .or_else(|| {
1128                    spec.capabilities
1129                        .adapters
1130                        .pty_sidecar
1131                        .as_ref()
1132                        .filter(|binding| *binding == &source.binding)
1133                })
1134        }
1135        AdapterFamily::Acp => spec
1136            .capabilities
1137            .transports
1138            .acp
1139            .as_ref()
1140            .map(|transport| &transport.adapter),
1141        AdapterFamily::OneShot => spec.capabilities.adapters.one_shot.as_ref(),
1142        AdapterFamily::Hook => spec.capabilities.adapters.hook.as_ref(),
1143        AdapterFamily::History
1144        | AdapterFamily::Resume
1145        | AdapterFamily::SessionOptions
1146        | AdapterFamily::CapabilityProbe
1147        | AdapterFamily::ManagedHook => None,
1148    }
1149}
1150
1151impl Default for Gate4AgentKernel {
1152    fn default() -> Self {
1153        Self::with_builtin_catalog()
1154    }
1155}
1156
1157#[cfg(test)]
1158mod tests {
1159    use super::*;
1160    use gate4agent_tool_protocol::{
1161        CancellationDisposition, CapabilityClass, CapabilityDescriptor, CapabilityObservation,
1162        CapabilityObservationEnvelope, CapabilityOwner, CapabilityRequestId,
1163        CapabilityRequestInput, CapabilityResult, CapabilityResultDelivery,
1164        CapabilityResultMetadata, CapabilityTerminalOutcome, ConsumerBoundCapabilityRequest,
1165        ConsumerId, GrantMode, PolicyDenial, PolicyGrant, PolicyKey, ResourceScopeId, ToolActorId,
1166        ToolAuthorityCommand, ToolCapabilityId, ToolProviderId,
1167    };
1168    use gate4agent_types::{
1169        AgentInstanceId, ApprovalLevel, CapabilityProbeRequest, ControlObservation, HistoryQuery,
1170        ObservationEnvelope, ProviderActivity, ProviderEvent, ProviderRuntimePolicy,
1171        ProviderSource, ResumeLaunchRequest, ResumeTarget, SessionGeneration,
1172        SessionOptionSelection, SessionStatus, StartRequest, TerminalSize, TransportKind,
1173    };
1174
1175    fn instance() -> AgentInstanceId {
1176        AgentInstanceId(11)
1177    }
1178
1179    fn verified_runtime_policy() -> ProviderRuntimePolicy {
1180        ProviderRuntimePolicy::new(true, true, true, true, true, true).unwrap()
1181    }
1182
1183    fn command(id: u64, command: ControlCommand) -> CommandEnvelope {
1184        CommandEnvelope {
1185            id: CommandId(id),
1186            command,
1187        }
1188    }
1189
1190    fn register(id: u64, agent: &str) -> CommandEnvelope {
1191        command(
1192            id,
1193            ControlCommand::Register {
1194                instance_id: instance(),
1195                agent_id: AgentId::new(agent).unwrap(),
1196                transport: TransportKind::Pty,
1197            },
1198        )
1199    }
1200
1201    fn tool_consumer() -> ConsumerId {
1202        ConsumerId::new("kernel-test-consumer").unwrap()
1203    }
1204
1205    fn tool_actor() -> ToolActorId {
1206        ToolActorId::new("kernel-test-actor").unwrap()
1207    }
1208
1209    fn tool_provider_id() -> ToolProviderId {
1210        ToolProviderId::new("kernel-browser-provider").unwrap()
1211    }
1212
1213    fn tool_capability_id() -> ToolCapabilityId {
1214        ToolCapabilityId::new("browser.snapshot").unwrap()
1215    }
1216
1217    fn tool_resource_scope() -> ResourceScopeId {
1218        ResourceScopeId::new("active-page").unwrap()
1219    }
1220
1221    /// `cursor`, `amp`, `no-resume-fixture`, and `pty-sidecar-fixture` are
1222    /// not part of the current fleet's built-in registry. Kernel
1223    /// adapter-family gating
1224    /// tests need provider capability shapes the four-member fleet does not
1225    /// naturally offer on its own -- e.g. a Pipe transport still resolving
1226    /// through the legacy `OneShotText` one-shot path (every fleet member's
1227    /// own Pipe transport is `StructuredJsonl`), a provider missing its
1228    /// History adapter, a provider with no Resume adapter, or a bare PTY
1229    /// sidecar binding -- so this clones a real fleet spec (`codex`, which
1230    /// carries PTY, History, and Resume) as the base and gives it a fixture
1231    /// identity, leaving every adapter binding it inherits pointed at a
1232    /// real, globally-registered implementation.
1233    fn legacy_fixture(id: &str) -> gate4agent_types::AgentSpec {
1234        let mut spec = builtin_registry().get_by_id("codex").unwrap().clone();
1235        spec.id = AgentId::new(id).unwrap();
1236        spec.detection.command = id.to_owned();
1237        spec.detection.aliases = Vec::new();
1238        spec.launch.program = id.to_owned();
1239        spec.expected_processes = vec![gate4agent_types::ProcessMatcher::Exact {
1240            name: id.to_owned(),
1241        }];
1242        spec
1243    }
1244
1245    /// `cursor`, standing in as the one legacy fixture that is both
1246    /// one-shot by Pipe and carries a session-option adapter -- bound to
1247    /// `claude`'s real one-shot implementation, the same resolver the
1248    /// kernel itself calls.
1249    fn legacy_one_shot_fixture_catalog() -> AgentRegistry {
1250        let claude = builtin_registry().get_by_id("claude").unwrap();
1251        let one_shot = claude.capabilities.adapters.one_shot.clone().unwrap();
1252        let session_options = claude.capabilities.adapters.session_options.clone();
1253        let mut cursor = legacy_fixture("cursor");
1254        cursor.capabilities.adapters.one_shot = Some(one_shot.clone());
1255        cursor.capabilities.adapters.session_options = session_options;
1256        cursor.capabilities.transports.pipe = Some(gate4agent_types::PipeTransportSpec {
1257            adapter: one_shot,
1258            protocol: PipeProtocol::OneShotText,
1259            launch_override: None,
1260            prompt_delivery: gate4agent_types::PipePromptDelivery::None,
1261        });
1262        AgentRegistry::new(builtin_registry().iter().cloned().chain([cursor])).unwrap()
1263    }
1264
1265    /// `pty-sidecar-fixture`, carrying a PTY sidecar binding reusing
1266    /// `claude`'s real, globally-registered Pipe binding -- the kernel's
1267    /// ingress gate only compares bindings for exact equality, so which
1268    /// real family member it is borrowed from is not load-bearing.
1269    fn legacy_pty_sidecar_fixture_catalog() -> AgentRegistry {
1270        let sidecar = builtin_registry()
1271            .get_by_id("claude")
1272            .unwrap()
1273            .capabilities
1274            .transports
1275            .pipe
1276            .clone()
1277            .unwrap()
1278            .adapter;
1279        let mut sidecar_fixture = legacy_fixture("pty-sidecar-fixture");
1280        // Cleared, not just left inherited from the `codex` skeleton: the
1281        // ingress-exactness assertion below reuses `codex`'s own Pipe
1282        // binding as the "foreign" one, and `declared_provider_binding`
1283        // accepts a Pipe-family source matching EITHER the sidecar binding
1284        // or the transport's own adapter, so leaving this set would let the
1285        // "foreign" binding match right back through it.
1286        sidecar_fixture.capabilities.transports.pipe = None;
1287        sidecar_fixture.capabilities.adapters.pty_sidecar = Some(sidecar);
1288        AgentRegistry::new(builtin_registry().iter().cloned().chain([sidecar_fixture])).unwrap()
1289    }
1290
1291    /// `amp`, carrying no History adapter.
1292    fn legacy_no_history_fixture_catalog() -> AgentRegistry {
1293        let mut amp = legacy_fixture("amp");
1294        amp.capabilities.adapters.history = None;
1295        AgentRegistry::new(builtin_registry().iter().cloned().chain([amp])).unwrap()
1296    }
1297
1298    /// `no-resume-fixture`, carrying every adapter family except Resume.
1299    fn legacy_no_resume_fixture_catalog() -> AgentRegistry {
1300        let mut no_resume_fixture = legacy_fixture("no-resume-fixture");
1301        no_resume_fixture.capabilities.adapters.resume = None;
1302        AgentRegistry::new(builtin_registry().iter().cloned().chain([no_resume_fixture])).unwrap()
1303    }
1304
1305    fn tool_provider() -> CapabilityProviderDescriptor {
1306        CapabilityProviderDescriptor {
1307            id: tool_provider_id(),
1308            owner: CapabilityOwner::Gate,
1309            capabilities: vec![CapabilityDescriptor::new(
1310                tool_capability_id(),
1311                CapabilityClass::Browser,
1312                "Return active page metadata",
1313            )
1314            .unwrap()],
1315        }
1316    }
1317
1318    fn other_tool_provider_id() -> ToolProviderId {
1319        ToolProviderId::new("kernel-browser-provider-secondary").unwrap()
1320    }
1321
1322    fn other_tool_provider() -> CapabilityProviderDescriptor {
1323        CapabilityProviderDescriptor {
1324            id: other_tool_provider_id(),
1325            owner: CapabilityOwner::Gate,
1326            capabilities: vec![CapabilityDescriptor::new(
1327                ToolCapabilityId::new("browser.snapshot.secondary").unwrap(),
1328                CapabilityClass::Browser,
1329                "Return secondary page metadata",
1330            )
1331            .unwrap()],
1332        }
1333    }
1334
1335    fn tool_request(
1336        local_id: u64,
1337        generation: SessionGeneration,
1338    ) -> ConsumerBoundCapabilityRequest {
1339        ConsumerBoundCapabilityRequest::new(
1340            tool_consumer(),
1341            tool_actor(),
1342            CapabilityRequestInput {
1343                local_id: CapabilityRequestId(local_id),
1344                instance_id: instance(),
1345                generation,
1346                provider_id: tool_provider_id(),
1347                capability_id: tool_capability_id(),
1348                resource_scope_id: tool_resource_scope(),
1349                approval_summary: "Read active page metadata".to_owned(),
1350                deadline_tick: 100,
1351                payload: br#"{"scope":"active-page"}"#.to_vec(),
1352            },
1353        )
1354    }
1355
1356    fn provider_bound_tool_request(
1357        binding_id: Option<ProviderBindingId>,
1358        request: ConsumerBoundCapabilityRequest,
1359    ) -> ProviderBoundCapabilityRequest {
1360        ProviderBoundCapabilityRequest::new(binding_id, request)
1361    }
1362
1363    fn tool_policy_grant(generation: SessionGeneration) -> PolicyGrant {
1364        PolicyGrant {
1365            key: PolicyKey {
1366                consumer_id: tool_consumer(),
1367                actor_id: tool_actor(),
1368                instance_id: instance(),
1369                generation,
1370                provider_id: tool_provider_id(),
1371                capability_id: tool_capability_id(),
1372                resource_scope_id: tool_resource_scope(),
1373            },
1374            mode: GrantMode::Allow,
1375        }
1376    }
1377
1378    fn tool_grant(generation: SessionGeneration, sequence: u64) -> ToolAuthorityEnvelope {
1379        ToolAuthorityEnvelope {
1380            sequence,
1381            command: ToolAuthorityCommand::SetGrant {
1382                grant: tool_policy_grant(generation),
1383            },
1384        }
1385    }
1386
1387    fn provider_runtime(sequence: u64, command: ProviderRuntimeCommand) -> BackendIngress {
1388        BackendIngress::ToolProvider(ProviderRuntimeEnvelope {
1389            sequence,
1390            command,
1391        })
1392    }
1393
1394    fn attach_tool_provider(kernel: &mut Gate4AgentKernel, sequence: u64) -> ProviderBindingId {
1395        let binding_id = ProviderBindingId(sequence);
1396        let step = kernel.step_control_plane(
1397            [provider_runtime(
1398                sequence,
1399                ProviderRuntimeCommand::Attach {
1400                    binding_id,
1401                    provider_id: tool_provider_id(),
1402                },
1403            )],
1404            [],
1405        );
1406        assert!(matches!(
1407            &step.ingress_outcomes[0],
1408            BackendIngressOutcome::ToolProvider(ProviderRuntimeCommandOutcome {
1409                result: Ok(ProviderRuntimeTransition::Attached),
1410                ..
1411            })
1412        ));
1413        binding_id
1414    }
1415
1416    fn successful_observation(
1417        effect: &ProviderBoundCapabilityEffectEnvelope,
1418    ) -> CapabilityObservationEnvelope {
1419        CapabilityObservationEnvelope {
1420            operation_id: effect.effect.operation_id,
1421            request_key: effect.effect.request_key.clone(),
1422            instance_id: effect.effect.instance_id,
1423            generation: effect.effect.generation,
1424            provider_id: effect.effect.provider_id.clone(),
1425            observation: CapabilityObservation::Succeeded {
1426                result: CapabilityResult {
1427                    metadata: CapabilityResultMetadata {
1428                        byte_len: 2,
1429                        media_type: Some("application/json".to_owned()),
1430                        truncated: false,
1431                        redacted_summary: Some("provider result".to_owned()),
1432                    },
1433                    delivery: CapabilityResultDelivery::Inline {
1434                        bytes: b"{}".to_vec(),
1435                    },
1436                },
1437            },
1438        }
1439    }
1440
1441    fn start_running(kernel: &mut Gate4AgentKernel) -> SessionGeneration {
1442        let starting = kernel.step(
1443            [
1444                register(1, "claude"),
1445                command(
1446                    2,
1447                    ControlCommand::Start {
1448                        instance_id: instance(),
1449                        runtime_policy: verified_runtime_policy(),
1450                        request: StartRequest {
1451                            working_directory: ".".to_owned(),
1452                            terminal_size: TerminalSize {
1453                                rows: 24,
1454                                columns: 80,
1455                            },
1456                            initial_prompt: None,
1457                            session_options: None,
1458                            approval_level: ApprovalLevel::default(),
1459                        },
1460                    },
1461                ),
1462            ],
1463            [],
1464        );
1465        let spawn = starting.effects[0].clone();
1466        let running = kernel.step(
1467            [],
1468            [ObservationEnvelope {
1469                operation_id: Some(spawn.operation_id),
1470                instance_id: spawn.instance_id,
1471                generation: spawn.generation,
1472                observation: ControlObservation::Spawned {
1473                    process_id: Some(123),
1474                },
1475            }],
1476        );
1477        assert!(running.integration_errors.is_empty());
1478        assert_eq!(running.snapshot.sessions[0].status, SessionStatus::Running);
1479        running.snapshot.sessions[0].generation
1480    }
1481
1482    #[test]
1483    fn unknown_provider_is_rejected_before_engine_mutation() {
1484        let mut kernel = Gate4AgentKernel::default();
1485        let step = kernel.step([register(1, "unknown-agent")], []);
1486
1487        assert!(matches!(
1488            step.command_outcomes[0].result,
1489            Err(KernelCommandError::UnknownAgent { .. })
1490        ));
1491        assert!(step.snapshot.sessions.is_empty());
1492        assert!(step.effects.is_empty());
1493        assert!(matches!(
1494            step.events[0].event,
1495            gate4agent_types::ControlEventKind::CommandRejected { .. }
1496        ));
1497    }
1498
1499    #[test]
1500    fn session_options_require_a_declared_pty_catalog_and_cross_the_effect_boundary() {
1501        let mut kernel = Gate4AgentKernel::default();
1502        kernel.step([register(1, "claude")], []);
1503        // `fastMode` is deliberately absent: unlike `cursor`'s (composed
1504        // into the model string), claude's has no `launch` application at
1505        // all, so it never reaches `applied` regardless of what is
1506        // selected -- only `effort` (a `--effort` flag) fully round-trips.
1507        let selection = SessionOptionSelection::new("opus").with_value("effort", "high");
1508        let accepted = kernel.step(
1509            [command(
1510                2,
1511                ControlCommand::Start {
1512                    instance_id: instance(),
1513                    runtime_policy: verified_runtime_policy(),
1514                    request: StartRequest {
1515                        working_directory: ".".to_owned(),
1516                        terminal_size: TerminalSize {
1517                            rows: 24,
1518                            columns: 80,
1519                        },
1520                        initial_prompt: None,
1521                        session_options: Some(selection.clone()),
1522                        approval_level: ApprovalLevel::default(),
1523                    },
1524                },
1525            )],
1526            [],
1527        );
1528        assert_eq!(accepted.command_outcomes[0].result, Ok(()));
1529        assert!(matches!(
1530            &accepted.effects[0].effect,
1531            gate4agent_types::ControlEffect::Spawn { request, .. }
1532                if request.session_options.as_ref() == Some(&selection)
1533        ));
1534
1535        let mut unsupported = Gate4AgentKernel::default();
1536        unsupported.step([register(1, "kimi")], []);
1537        let rejected = unsupported.step(
1538            [command(
1539                2,
1540                ControlCommand::Start {
1541                    instance_id: instance(),
1542                    runtime_policy: verified_runtime_policy(),
1543                    request: StartRequest {
1544                        working_directory: ".".to_owned(),
1545                        terminal_size: TerminalSize {
1546                            rows: 24,
1547                            columns: 80,
1548                        },
1549                        initial_prompt: None,
1550                        session_options: Some(SessionOptionSelection::new("opus")),
1551                        approval_level: ApprovalLevel::default(),
1552                    },
1553                },
1554            )],
1555            [],
1556        );
1557        assert!(matches!(
1558            rejected.command_outcomes[0].result,
1559            Err(KernelCommandError::UnsupportedCapability {
1560                capability: "pty-session-options",
1561                ..
1562            })
1563        ));
1564        assert!(rejected.effects.is_empty());
1565
1566        let mut expanded = Gate4AgentKernel::default();
1567        expanded.step([register(1, "claude")], []);
1568        let started = expanded.step(
1569            [command(
1570                2,
1571                ControlCommand::Start {
1572                    instance_id: instance(),
1573                    runtime_policy: verified_runtime_policy(),
1574                    request: StartRequest {
1575                        working_directory: ".".to_owned(),
1576                        terminal_size: TerminalSize {
1577                            rows: 24,
1578                            columns: 80,
1579                        },
1580                        initial_prompt: None,
1581                        session_options: Some(SessionOptionSelection::new("opus")),
1582                        approval_level: ApprovalLevel::default(),
1583                    },
1584                },
1585            )],
1586            [],
1587        );
1588        let expected = SessionOptionSelection::new("opus").with_value("effort", "high");
1589        assert_eq!(
1590            started.snapshot.sessions[0].session_options.as_ref(),
1591            Some(&expected)
1592        );
1593        assert!(matches!(
1594            &started.effects[0].effect,
1595            gate4agent_types::ControlEffect::Spawn { request, .. }
1596                if request.session_options.as_ref() == Some(&expected)
1597        ));
1598
1599        let mut pipe = Gate4AgentKernel::default();
1600        pipe.step(
1601            [command(
1602                1,
1603                ControlCommand::Register {
1604                    instance_id: instance(),
1605                    agent_id: AgentId::new("codex").unwrap(),
1606                    transport: TransportKind::Pipe,
1607                },
1608            )],
1609            [],
1610        );
1611        let rejected = pipe.step(
1612            [command(
1613                2,
1614                ControlCommand::Start {
1615                    instance_id: instance(),
1616                    runtime_policy: verified_runtime_policy(),
1617                    request: StartRequest {
1618                        working_directory: ".".to_owned(),
1619                        terminal_size: TerminalSize {
1620                            rows: 24,
1621                            columns: 80,
1622                        },
1623                        initial_prompt: Some("hello".to_owned()),
1624                        session_options: Some(SessionOptionSelection::new("gpt-5.5")),
1625                        approval_level: ApprovalLevel::default(),
1626                    },
1627                },
1628            )],
1629            [],
1630        );
1631        assert!(matches!(
1632            rejected.command_outcomes[0].result,
1633            Err(KernelCommandError::UnsupportedCapability {
1634                capability: "pty-session-options",
1635                ..
1636            })
1637        ));
1638    }
1639
1640    /// The one-shot pipe path fills in a provider's own session-option
1641    /// defaults when a `Start` names none, and the same resolved selection
1642    /// reaches the spawn effect.
1643    ///
1644    /// Subject is `cursor`, and it has to be: this used to drive the case
1645    /// with `claude`, which stopped being a one-shot provider when its pipe
1646    /// transport became `StructuredJsonl` (pinned, green, in
1647    /// `gate4agent-catalog`'s own builtin test). The kernel's one-shot
1648    /// branch is gated on `PipeProtocol::OneShotText`, so with `claude` it
1649    /// simply never ran and the test asserted defaults against a path it
1650    /// was no longer on -- red ever since, which is how a test stops being
1651    /// read. `cursor` is the one provider that is both one-shot by pipe and
1652    /// carries a session-option adapter, so it is the only subject that
1653    /// exercises what this test is named for.
1654    ///
1655    /// The expectation is derived from `resolve_one_shot_plan`, the same
1656    /// resolver the kernel calls, rather than frozen as a literal model
1657    /// name. A frozen literal is what rotted the previous version: the
1658    /// defaults belong to the catalog and move with it, while what this
1659    /// test is actually about -- that the kernel APPLIES them instead of
1660    /// leaving `session_options` at `None` -- does not.
1661    #[test]
1662    fn one_shot_pipe_defaults_and_validates_options_before_effect_creation() {
1663        let mut cursor = Gate4AgentKernel::new(legacy_one_shot_fixture_catalog());
1664        cursor.step(
1665            [command(
1666                1,
1667                ControlCommand::Register {
1668                    instance_id: instance(),
1669                    agent_id: AgentId::new("cursor").unwrap(),
1670                    transport: TransportKind::Pipe,
1671                },
1672            )],
1673            [],
1674        );
1675        let started = cursor.step(
1676            [command(
1677                2,
1678                ControlCommand::Start {
1679                    instance_id: instance(),
1680                    runtime_policy: verified_runtime_policy(),
1681                    request: StartRequest {
1682                        working_directory: ".".to_owned(),
1683                        terminal_size: TerminalSize {
1684                            rows: 24,
1685                            columns: 80,
1686                        },
1687                        initial_prompt: Some("summarize".to_owned()),
1688                        session_options: None,
1689                        approval_level: ApprovalLevel::default(),
1690                    },
1691                },
1692            )],
1693            [],
1694        );
1695        let spec = cursor
1696            .catalog()
1697            .get_by_id("cursor")
1698            .expect("cursor is a legacy fixture provider");
1699        let binding = spec
1700            .capabilities
1701            .adapters
1702            .one_shot
1703            .as_ref()
1704            .expect("cursor declares a one-shot adapter");
1705        let expected = resolve_one_shot_plan(&binding.id, &spec.launch, "summarize", None)
1706            .expect("cursor resolves its own defaults")
1707            .applied;
1708        // Stated on its own: the defect this guards against is the kernel
1709        // leaving the field untouched, and `assert_eq!` against a resolved
1710        // value would report that as a mismatch rather than as the absence
1711        // it is.
1712        assert!(started.snapshot.sessions[0].session_options.is_some());
1713        assert_eq!(
1714            started.snapshot.sessions[0].session_options.as_ref(),
1715            Some(&expected)
1716        );
1717        assert!(matches!(
1718            &started.effects[0].effect,
1719            gate4agent_types::ControlEffect::Spawn {
1720                transport: TransportKind::Pipe,
1721                request,
1722                ..
1723            } if request.session_options.as_ref() == Some(&expected)
1724        ));
1725
1726        let mut amp = Gate4AgentKernel::new(legacy_one_shot_fixture_catalog());
1727        amp.step(
1728            [command(
1729                3,
1730                ControlCommand::Register {
1731                    instance_id: instance(),
1732                    agent_id: AgentId::new("cursor").unwrap(),
1733                    transport: TransportKind::Pipe,
1734                },
1735            )],
1736            [],
1737        );
1738        let rejected = amp.step(
1739            [command(
1740                4,
1741                ControlCommand::Start {
1742                    instance_id: instance(),
1743                    runtime_policy: verified_runtime_policy(),
1744                    request: StartRequest {
1745                        working_directory: ".".to_owned(),
1746                        terminal_size: TerminalSize {
1747                            rows: 24,
1748                            columns: 80,
1749                        },
1750                        initial_prompt: Some("summarize".to_owned()),
1751                        session_options: Some(SessionOptionSelection::new("unknown-model")),
1752                        approval_level: ApprovalLevel::default(),
1753                    },
1754                },
1755            )],
1756            [],
1757        );
1758        assert!(matches!(
1759            rejected.command_outcomes[0].result,
1760            Err(KernelCommandError::InvalidSessionOptions { .. })
1761        ));
1762        assert!(rejected.effects.is_empty());
1763
1764        // `codex` drove this half until its pipe transport became
1765        // `StructuredJsonl` too, which routes it away from the one-shot
1766        // prompt validation the half exists to check.
1767        let mut missing_prompt = Gate4AgentKernel::new(legacy_one_shot_fixture_catalog());
1768        missing_prompt.step(
1769            [command(
1770                5,
1771                ControlCommand::Register {
1772                    instance_id: instance(),
1773                    agent_id: AgentId::new("cursor").unwrap(),
1774                    transport: TransportKind::Pipe,
1775                },
1776            )],
1777            [],
1778        );
1779        let rejected = missing_prompt.step(
1780            [command(
1781                6,
1782                ControlCommand::Start {
1783                    instance_id: instance(),
1784                    runtime_policy: verified_runtime_policy(),
1785                    request: StartRequest {
1786                        working_directory: ".".to_owned(),
1787                        terminal_size: TerminalSize {
1788                            rows: 24,
1789                            columns: 80,
1790                        },
1791                        initial_prompt: None,
1792                        session_options: None,
1793                        approval_level: ApprovalLevel::default(),
1794                    },
1795                },
1796            )],
1797            [],
1798        );
1799        assert!(matches!(
1800            rejected.command_outcomes[0].result,
1801            Err(KernelCommandError::InvalidSessionOptions { .. })
1802        ));
1803        assert!(rejected.effects.is_empty());
1804    }
1805
1806    /// Capability-probe existed to serve `cursor`'s `--list-models`, and
1807    /// `cursor` is not part of the fleet: no fleet provider declares this
1808    /// adapter, so the kernel rejects it before creating an effect for all
1809    /// four; and a provider that DOES declare a capability-probe binding
1810    /// still fails closed at resolution -- the binding is only valid
1811    /// against a consumer-extended registry (built locally, standing in for
1812    /// what the wider reference catalog used to provide), never the
1813    /// process-global one the kernel resolves against.
1814    #[test]
1815    fn capability_probe_is_unsupported_fleet_wide_and_fails_closed_on_an_unavailable_binding() {
1816        for id in ["claude", "codex", "grok", "kimi"] {
1817            let mut kernel = Gate4AgentKernel::default();
1818            kernel.step([register(1, id)], []);
1819            let rejected = kernel.step(
1820                [command(
1821                    2,
1822                    ControlCommand::ProbeCapabilities {
1823                        instance_id: instance(),
1824                        request: CapabilityProbeRequest {
1825                            working_directory: ".".to_owned(),
1826                        },
1827                    },
1828                )],
1829                [],
1830            );
1831            assert!(
1832                matches!(
1833                    rejected.command_outcomes[0].result,
1834                    Err(KernelCommandError::UnsupportedCapability {
1835                        capability: "capability-probe",
1836                        ..
1837                    })
1838                ),
1839                "{id}"
1840            );
1841            assert!(rejected.effects.is_empty(), "{id}");
1842        }
1843
1844        let probe_binding = gate4agent_types::AdapterBinding::new(
1845            gate4agent_types::AdapterId::new("cursor").unwrap(),
1846            "cursor-capability-probe/v1",
1847            gate4agent_types::AdapterVerification::Reference,
1848        )
1849        .unwrap();
1850        let local_adapters = gate4agent_catalog::AdapterRegistry::new(
1851            gate4agent_catalog::builtin_adapter_registry()
1852                .iter()
1853                .cloned()
1854                .chain([gate4agent_catalog::AdapterDescriptor {
1855                    family: AdapterFamily::CapabilityProbe,
1856                    binding: probe_binding.clone(),
1857                    agents: vec![AgentId::new("cursor").unwrap()],
1858                }]),
1859        )
1860        .unwrap();
1861        let mut cursor = legacy_fixture("cursor");
1862        cursor.capabilities.adapters.capability_probe = Some(probe_binding);
1863        let catalog = AgentRegistry::new_with_adapters(
1864            builtin_registry().iter().cloned().chain([cursor]),
1865            &local_adapters,
1866        )
1867        .unwrap();
1868        let mut kernel = Gate4AgentKernel::new(catalog);
1869        kernel.step([register(1, "cursor")], []);
1870        let rejected = kernel.step(
1871            [command(
1872                2,
1873                ControlCommand::ProbeCapabilities {
1874                    instance_id: instance(),
1875                    request: CapabilityProbeRequest {
1876                        working_directory: ".".to_owned(),
1877                    },
1878                },
1879            )],
1880            [],
1881        );
1882        assert!(matches!(
1883            rejected.command_outcomes[0].result,
1884            Err(KernelCommandError::InvalidCapabilityProbe { .. })
1885        ));
1886        assert!(rejected.effects.is_empty());
1887    }
1888
1889    #[test]
1890    fn external_provider_ingress_requires_the_declared_family_binding() {
1891        let mut kernel = Gate4AgentKernel::default();
1892        kernel.step([register(1, "grok")], []);
1893        let started = kernel.step(
1894            [command(
1895                2,
1896                ControlCommand::Start {
1897                    instance_id: instance(),
1898                    runtime_policy: verified_runtime_policy(),
1899                    request: StartRequest {
1900                        working_directory: ".".to_owned(),
1901                        terminal_size: TerminalSize {
1902                            rows: 24,
1903                            columns: 80,
1904                        },
1905                        initial_prompt: None,
1906                        session_options: None,
1907                        approval_level: ApprovalLevel::default(),
1908                    },
1909                },
1910            )],
1911            [],
1912        );
1913        let generation = started.snapshot.sessions[0].generation;
1914        // `AdapterFamily::Hook` is retired (owner ruling 2026-09-25) and
1915        // `declared_provider_binding` never resolves a binding for it
1916        // anymore. `Acp` is the only ingress-checked family grok declares
1917        // (it has no PtySemantic/Pipe/OneShot adapter), and it exercises the
1918        // same declared-binding validation this test is about.
1919        let grok_acp = kernel
1920            .catalog()
1921            .get_by_id("grok")
1922            .unwrap()
1923            .capabilities
1924            .transports
1925            .acp
1926            .clone()
1927            .unwrap()
1928            .adapter;
1929        let accepted = kernel.step(
1930            [command(
1931                3,
1932                ControlCommand::IngestProvider {
1933                    instance_id: instance(),
1934                    generation,
1935                    source: ProviderSource {
1936                        family: AdapterFamily::Acp,
1937                        binding: grok_acp,
1938                    },
1939                    source_sequence: 1,
1940                    events: vec![ProviderEvent::TurnStarted {
1941                        prompt: Some("ground external".to_owned()),
1942                    }],
1943                },
1944            )],
1945            [],
1946        );
1947        assert_eq!(accepted.command_outcomes[0].result, Ok(()));
1948        assert_eq!(
1949            accepted.snapshot.sessions[0].provider.activity,
1950            ProviderActivity::Working
1951        );
1952
1953        let kimi_acp = kernel
1954            .catalog()
1955            .get_by_id("kimi")
1956            .unwrap()
1957            .capabilities
1958            .transports
1959            .acp
1960            .clone()
1961            .unwrap()
1962            .adapter;
1963        let rejected = kernel.step(
1964            [command(
1965                4,
1966                ControlCommand::IngestProvider {
1967                    instance_id: instance(),
1968                    generation,
1969                    source: ProviderSource {
1970                        family: AdapterFamily::Acp,
1971                        binding: kimi_acp,
1972                    },
1973                    source_sequence: 2,
1974                    events: vec![ProviderEvent::Ready],
1975                },
1976            )],
1977            [],
1978        );
1979        assert!(matches!(
1980            rejected.command_outcomes[0].result,
1981            Err(KernelCommandError::InvalidProviderSource { .. })
1982        ));
1983    }
1984
1985    #[test]
1986    fn pipe_provider_ingress_accepts_exact_pty_sidecar_binding_only() {
1987        let mut kernel = Gate4AgentKernel::new(legacy_pty_sidecar_fixture_catalog());
1988        kernel.step([register(1, "pty-sidecar-fixture")], []);
1989        let started = kernel.step(
1990            [command(
1991                2,
1992                ControlCommand::Start {
1993                    instance_id: instance(),
1994                    runtime_policy: verified_runtime_policy(),
1995                    request: StartRequest {
1996                        working_directory: ".".to_owned(),
1997                        terminal_size: TerminalSize {
1998                            rows: 24,
1999                            columns: 80,
2000                        },
2001                        initial_prompt: None,
2002                        session_options: None,
2003                        approval_level: ApprovalLevel::default(),
2004                    },
2005                },
2006            )],
2007            [],
2008        );
2009        let generation = started.snapshot.sessions[0].generation;
2010        let sidecar = kernel
2011            .catalog()
2012            .get_by_id("pty-sidecar-fixture")
2013            .unwrap()
2014            .capabilities
2015            .adapters
2016            .pty_sidecar
2017            .clone()
2018            .unwrap();
2019        let accepted = kernel.step(
2020            [command(
2021                3,
2022                ControlCommand::IngestProvider {
2023                    instance_id: instance(),
2024                    generation,
2025                    source: ProviderSource {
2026                        family: AdapterFamily::Pipe,
2027                        binding: sidecar,
2028                    },
2029                    source_sequence: 1,
2030                    events: vec![ProviderEvent::Ready],
2031                },
2032            )],
2033            [],
2034        );
2035        assert_eq!(accepted.command_outcomes[0].result, Ok(()));
2036
2037        let foreign_pipe = kernel
2038            .catalog()
2039            .get_by_id("codex")
2040            .unwrap()
2041            .capabilities
2042            .transports
2043            .pipe
2044            .as_ref()
2045            .unwrap()
2046            .adapter
2047            .clone();
2048        let rejected = kernel.step(
2049            [command(
2050                4,
2051                ControlCommand::IngestProvider {
2052                    instance_id: instance(),
2053                    generation,
2054                    source: ProviderSource {
2055                        family: AdapterFamily::Pipe,
2056                        binding: foreign_pipe,
2057                    },
2058                    source_sequence: 2,
2059                    events: vec![ProviderEvent::Ready],
2060                },
2061            )],
2062            [],
2063        );
2064        assert!(matches!(
2065            rejected.command_outcomes[0].result,
2066            Err(KernelCommandError::InvalidProviderSource { .. })
2067        ));
2068    }
2069
2070    #[test]
2071    fn unsupported_provider_transport_is_rejected_before_registration() {
2072        let mut kernel = Gate4AgentKernel::default();
2073        let outcome = kernel.step(
2074            [command(
2075                1,
2076                ControlCommand::Register {
2077                    instance_id: instance(),
2078                    agent_id: AgentId::new("grok").unwrap(),
2079                    transport: TransportKind::Pipe,
2080                },
2081            )],
2082            [],
2083        );
2084        assert!(matches!(
2085            &outcome.command_outcomes[0].result,
2086            Err(KernelCommandError::UnsupportedTransport {
2087                transport: TransportKind::Pipe,
2088                ..
2089            })
2090        ));
2091        assert!(outcome.snapshot.sessions.is_empty());
2092    }
2093
2094    #[test]
2095    fn history_commands_require_a_declared_adapter_before_effect_creation() {
2096        let mut supported = Gate4AgentKernel::default();
2097        supported.step([register(1, "grok")], []);
2098        let accepted = supported.step(
2099            [command(
2100                2,
2101                ControlCommand::DiscoverHistory {
2102                    instance_id: instance(),
2103                    query: HistoryQuery {
2104                        working_directory: None,
2105                        limit: 8,
2106                    },
2107                },
2108            )],
2109            [],
2110        );
2111        assert_eq!(accepted.command_outcomes[0].result, Ok(()));
2112        assert!(matches!(
2113            accepted.effects[0].effect,
2114            gate4agent_types::ControlEffect::DiscoverHistory { .. }
2115        ));
2116
2117        // Subject is a legacy `amp` fixture with no history adapter. Every
2118        // fleet member now declares one (pinned in `gate4agent-adapters`'
2119        // own registry test), so this shape does not occur naturally in the
2120        // fleet anymore and needs a constructed fixture to stay exercised at
2121        // all.
2122        let mut unsupported = Gate4AgentKernel::new(legacy_no_history_fixture_catalog());
2123        unsupported.step([register(1, "amp")], []);
2124        let rejected = unsupported.step(
2125            [command(
2126                2,
2127                ControlCommand::DiscoverHistory {
2128                    instance_id: instance(),
2129                    query: HistoryQuery {
2130                        working_directory: None,
2131                        limit: 8,
2132                    },
2133                },
2134            )],
2135            [],
2136        );
2137        assert!(matches!(
2138            rejected.command_outcomes[0].result,
2139            Err(KernelCommandError::UnsupportedCapability {
2140                capability: "history",
2141                ..
2142            })
2143        ));
2144        assert!(rejected.effects.is_empty());
2145    }
2146
2147    #[test]
2148    fn resume_requires_a_declared_adapter_and_supported_transport_before_engine_mutation() {
2149        let resume = |initial_prompt| {
2150            command(
2151                2,
2152                ControlCommand::Resume {
2153                    instance_id: instance(),
2154                    target: ResumeTarget::CurrentProvider,
2155                    runtime_policy: verified_runtime_policy(),
2156                    request: ResumeLaunchRequest {
2157                        working_directory: ".".to_owned(),
2158                        terminal_size: TerminalSize {
2159                            rows: 24,
2160                            columns: 80,
2161                        },
2162                        initial_prompt,
2163                    },
2164                },
2165            )
2166        };
2167
2168        let mut unsupported = Gate4AgentKernel::new(legacy_no_resume_fixture_catalog());
2169        unsupported.step([register(1, "no-resume-fixture")], []);
2170        let rejected = unsupported.step([resume(None)], []);
2171        assert!(matches!(
2172            rejected.command_outcomes[0].result,
2173            Err(KernelCommandError::UnsupportedCapability {
2174                capability: "resume",
2175                ..
2176            })
2177        ));
2178        assert!(rejected.effects.is_empty());
2179
2180        let mut wrong_transport = Gate4AgentKernel::default();
2181        wrong_transport.step(
2182            [command(
2183                1,
2184                ControlCommand::Register {
2185                    instance_id: instance(),
2186                    agent_id: AgentId::new("codex").unwrap(),
2187                    transport: TransportKind::Pipe,
2188                },
2189            )],
2190            [],
2191        );
2192        let rejected = wrong_transport.step([resume(Some("continue".to_owned()))], []);
2193        assert!(matches!(
2194            rejected.command_outcomes[0].result,
2195            Err(KernelCommandError::Control(
2196                ControlError::MissingProviderSession
2197            ))
2198        ));
2199        assert!(rejected.effects.is_empty());
2200    }
2201
2202    #[test]
2203    fn undeclared_agent_commands_are_rejected_before_effect_creation() {
2204        let mut kernel = Gate4AgentKernel::default();
2205        let started = kernel.step(
2206            [
2207                register(1, "grok"),
2208                command(
2209                    2,
2210                    ControlCommand::Start {
2211                        instance_id: instance(),
2212                        runtime_policy: verified_runtime_policy(),
2213                        request: StartRequest {
2214                            working_directory: ".".to_owned(),
2215                            terminal_size: TerminalSize {
2216                                rows: 24,
2217                                columns: 80,
2218                            },
2219                            initial_prompt: None,
2220                            session_options: None,
2221                            approval_level: ApprovalLevel::default(),
2222                        },
2223                    },
2224                ),
2225            ],
2226            [],
2227        );
2228        let spawn = started.effects[0].clone();
2229        kernel.step(
2230            [],
2231            [ObservationEnvelope {
2232                operation_id: Some(spawn.operation_id),
2233                instance_id: spawn.instance_id,
2234                generation: spawn.generation,
2235                observation: ControlObservation::Spawned {
2236                    process_id: Some(123),
2237                },
2238            }],
2239        );
2240
2241        let rejected = kernel.step(
2242            [command(
2243                3,
2244                ControlCommand::SendInput {
2245                    instance_id: instance(),
2246                    action: InputAction::AgentCommand(gate4agent_types::AgentCommand {
2247                        agent_id: AgentId::new("grok").unwrap(),
2248                        name: "help".to_owned(),
2249                        arguments: Vec::new(),
2250                    }),
2251                },
2252            )],
2253            [],
2254        );
2255        assert!(matches!(
2256            rejected.command_outcomes[0].result,
2257            Err(KernelCommandError::UnsupportedCapability {
2258                capability: "agent-commands",
2259                ..
2260            })
2261        ));
2262        assert!(rejected.effects.is_empty());
2263    }
2264
2265    #[test]
2266    fn command_phase_precedes_observation_phase() {
2267        let mut kernel = Gate4AgentKernel::default();
2268        let first = kernel.step(
2269            [
2270                register(1, "claude"),
2271                command(
2272                    2,
2273                    ControlCommand::Start {
2274                        instance_id: instance(),
2275                        runtime_policy: verified_runtime_policy(),
2276                        request: StartRequest {
2277                            working_directory: ".".to_owned(),
2278                            terminal_size: TerminalSize {
2279                                rows: 24,
2280                                columns: 80,
2281                            },
2282                            initial_prompt: None,
2283                            session_options: None,
2284                            approval_level: ApprovalLevel::default(),
2285                        },
2286                    },
2287                ),
2288            ],
2289            [],
2290        );
2291        let spawn = first.effects[0].clone();
2292        let running = kernel.step(
2293            [],
2294            [ObservationEnvelope {
2295                operation_id: Some(spawn.operation_id),
2296                instance_id: spawn.instance_id,
2297                generation: spawn.generation,
2298                observation: ControlObservation::Spawned {
2299                    process_id: Some(123),
2300                },
2301            }],
2302        );
2303        assert_eq!(running.snapshot.sessions[0].status, SessionStatus::Running);
2304
2305        let raced = kernel.step(
2306            [command(
2307                3,
2308                ControlCommand::Stop {
2309                    instance_id: instance(),
2310                    force: false,
2311                },
2312            )],
2313            [ObservationEnvelope {
2314                operation_id: None,
2315                instance_id: instance(),
2316                generation: spawn.generation,
2317                observation: ControlObservation::ProcessExited {
2318                    exit_code: Some(0),
2319                    final_terminal: None,
2320                },
2321            }],
2322        );
2323
2324        assert!(raced.command_outcomes[0].result.is_ok());
2325        assert!(raced.effects.is_empty());
2326        assert_eq!(
2327            raced.snapshot.sessions[0].status,
2328            SessionStatus::Exited { exit_code: Some(0) }
2329        );
2330    }
2331
2332    #[test]
2333    fn identical_batches_produce_identical_step() {
2334        fn run() -> KernelStep {
2335            let mut kernel = Gate4AgentKernel::default();
2336            kernel.step(
2337                [
2338                    register(1, "claude"),
2339                    command(
2340                        2,
2341                        ControlCommand::Start {
2342                            instance_id: instance(),
2343                            runtime_policy: verified_runtime_policy(),
2344                            request: StartRequest {
2345                                working_directory: ".".to_owned(),
2346                                terminal_size: TerminalSize {
2347                                    rows: 24,
2348                                    columns: 80,
2349                                },
2350                                initial_prompt: None,
2351                                session_options: None,
2352                                approval_level: ApprovalLevel::default(),
2353                            },
2354                        },
2355                    ),
2356                ],
2357                [],
2358            )
2359        }
2360
2361        assert_eq!(run(), run());
2362    }
2363
2364    #[test]
2365    fn non_running_control_sessions_are_inactive_in_the_tool_engine() {
2366        let mut kernel =
2367            Gate4AgentKernel::with_tool_providers(builtin_registry().clone(), [tool_provider()])
2368                .unwrap();
2369        let registered = kernel.step([register(1, "claude")], []);
2370        assert_eq!(
2371            registered.backend_snapshot.tools.instance_states,
2372            vec![(instance(), ToolInstanceState::Inactive)]
2373        );
2374
2375        let starting = kernel.step(
2376            [command(
2377                2,
2378                ControlCommand::Start {
2379                    instance_id: instance(),
2380                    runtime_policy: verified_runtime_policy(),
2381                    request: StartRequest {
2382                        working_directory: ".".to_owned(),
2383                        terminal_size: TerminalSize {
2384                            rows: 24,
2385                            columns: 80,
2386                        },
2387                        initial_prompt: None,
2388                        session_options: None,
2389                        approval_level: ApprovalLevel::default(),
2390                    },
2391                },
2392            )],
2393            [],
2394        );
2395        assert_eq!(
2396            starting.backend_snapshot.tools.instance_states,
2397            vec![(instance(), ToolInstanceState::Inactive)]
2398        );
2399    }
2400
2401    #[test]
2402    fn tool_ingress_precedes_control_observation_activation() {
2403        let mut kernel =
2404            Gate4AgentKernel::with_tool_providers(builtin_registry().clone(), [tool_provider()])
2405                .unwrap();
2406        attach_tool_provider(&mut kernel, 1);
2407        let starting = kernel.step(
2408            [
2409                register(1, "claude"),
2410                command(
2411                    2,
2412                    ControlCommand::Start {
2413                        instance_id: instance(),
2414                        runtime_policy: verified_runtime_policy(),
2415                        request: StartRequest {
2416                            working_directory: ".".to_owned(),
2417                            terminal_size: TerminalSize {
2418                                rows: 24,
2419                                columns: 80,
2420                            },
2421                            initial_prompt: None,
2422                            session_options: None,
2423                            approval_level: ApprovalLevel::default(),
2424                        },
2425                    },
2426                ),
2427            ],
2428            [],
2429        );
2430        let spawn = starting.effects[0].clone();
2431        let step = kernel.step_control_plane(
2432            [BackendIngress::ToolRequest(provider_bound_tool_request(
2433                Some(ProviderBindingId(1)),
2434                tool_request(1, spawn.generation),
2435            ))],
2436            [ObservationEnvelope {
2437                operation_id: Some(spawn.operation_id),
2438                instance_id: spawn.instance_id,
2439                generation: spawn.generation,
2440                observation: ControlObservation::Spawned {
2441                    process_id: Some(123),
2442                },
2443            }],
2444        );
2445
2446        let BackendIngressOutcome::ToolRequest(outcome) = &step.ingress_outcomes[0] else {
2447            panic!("expected tool request outcome");
2448        };
2449        assert_eq!(
2450            outcome.result,
2451            Ok(PolicyDecision::Deny(PolicyDenial::InactiveInstance))
2452        );
2453        assert!(outcome.accepted_sequence.is_some());
2454        assert_eq!(
2455            step.backend_snapshot.tools.instance_states,
2456            vec![(instance(), ToolInstanceState::Active)]
2457        );
2458        assert!(step.tool_effects.is_empty());
2459        assert!(matches!(
2460            step.tool_completions.completions[0].outcome,
2461            CapabilityTerminalOutcome::PolicyDenied {
2462                reason: PolicyDenial::InactiveInstance
2463            }
2464        ));
2465    }
2466
2467    #[test]
2468    fn ordered_authority_and_request_outcomes_are_exactly_correlated() {
2469        let mut kernel =
2470            Gate4AgentKernel::with_tool_providers(builtin_registry().clone(), [tool_provider()])
2471                .unwrap();
2472        attach_tool_provider(&mut kernel, 1);
2473        let generation = start_running(&mut kernel);
2474        let request = tool_request(7, generation);
2475        let request_key = request.key();
2476        let step = kernel.step_control_plane(
2477            [
2478                BackendIngress::ToolAuthority(tool_grant(generation, 1)),
2479                BackendIngress::ToolRequest(provider_bound_tool_request(
2480                    Some(ProviderBindingId(1)),
2481                    request,
2482                )),
2483            ],
2484            [],
2485        );
2486
2487        assert!(matches!(
2488            &step.ingress_outcomes[0],
2489            BackendIngressOutcome::ToolAuthority(ToolAuthorityCommandOutcome {
2490                sequence: 1,
2491                result: Ok(ToolAuthorityOutcome::GrantSet),
2492            })
2493        ));
2494        let BackendIngressOutcome::ToolRequest(outcome) = &step.ingress_outcomes[1] else {
2495            panic!("expected tool request outcome");
2496        };
2497        assert_eq!(outcome.request_key, request_key);
2498        assert!(outcome.accepted_sequence.is_some());
2499        assert_eq!(outcome.result, Ok(PolicyDecision::Allow));
2500        assert_eq!(step.tool_effects.len(), 1);
2501        assert_eq!(step.tool_effects[0].effect.request_key, request_key);
2502
2503        let regressed = kernel.step_control_plane(
2504            [BackendIngress::ToolAuthority(tool_grant(generation, 1))],
2505            [],
2506        );
2507        assert!(matches!(
2508            &regressed.ingress_outcomes[0],
2509            BackendIngressOutcome::ToolAuthority(ToolAuthorityCommandOutcome {
2510                sequence: 1,
2511                result: Err(KernelToolError::Engine(
2512                    ToolEngineError::AuthoritySequenceRegressed {
2513                        current: 1,
2514                        requested: 1,
2515                    }
2516                )),
2517            })
2518        ));
2519    }
2520
2521    #[test]
2522    fn stop_fences_queued_tool_work_then_remove_register_advances_generation() {
2523        let mut kernel =
2524            Gate4AgentKernel::with_tool_providers(builtin_registry().clone(), [tool_provider()])
2525                .unwrap();
2526        attach_tool_provider(&mut kernel, 1);
2527        let generation = start_running(&mut kernel);
2528        let granted = kernel.step_control_plane(
2529            [BackendIngress::ToolAuthority(tool_grant(generation, 1))],
2530            [],
2531        );
2532        assert!(granted.integration_errors.is_empty());
2533        let request = tool_request(9, generation);
2534        let request_key = request.key();
2535
2536        let stopped = kernel.step_control_plane(
2537            [
2538                BackendIngress::ToolRequest(provider_bound_tool_request(
2539                    Some(ProviderBindingId(1)),
2540                    request,
2541                )),
2542                BackendIngress::Control(command(
2543                    3,
2544                    ControlCommand::Stop {
2545                        instance_id: instance(),
2546                        force: false,
2547                    },
2548                )),
2549            ],
2550            [],
2551        );
2552
2553        assert!(stopped.integration_errors.is_empty());
2554        assert!(stopped.tool_effects.is_empty());
2555        let completion = &stopped.tool_completions.completions[0];
2556        assert_eq!(completion.request_key, request_key);
2557        assert!(matches!(
2558            completion.outcome,
2559            CapabilityTerminalOutcome::InstanceClosed { .. }
2560        ));
2561
2562        let exited = kernel.step(
2563            [],
2564            [ObservationEnvelope {
2565                operation_id: None,
2566                instance_id: instance(),
2567                generation,
2568                observation: ControlObservation::ProcessExited {
2569                    exit_code: Some(0),
2570                    final_terminal: None,
2571                },
2572            }],
2573        );
2574        assert_eq!(
2575            exited.snapshot.sessions[0].status,
2576            SessionStatus::Exited { exit_code: Some(0) }
2577        );
2578
2579        let step = kernel.step_control_plane(
2580            [
2581                BackendIngress::Control(command(
2582                    4,
2583                    ControlCommand::Remove {
2584                        instance_id: instance(),
2585                    },
2586                )),
2587                BackendIngress::Control(register(5, "claude")),
2588            ],
2589            [],
2590        );
2591        assert!(step.integration_errors.is_empty());
2592        assert_eq!(step.snapshot.sessions.len(), 1);
2593        assert_eq!(step.snapshot.sessions[0].generation, SessionGeneration(2));
2594        assert_eq!(
2595            step.backend_snapshot.tools.generations,
2596            vec![(instance(), SessionGeneration(2))]
2597        );
2598        assert_eq!(
2599            step.backend_snapshot.tools.instance_states,
2600            vec![(instance(), ToolInstanceState::Inactive)]
2601        );
2602    }
2603
2604    #[test]
2605    fn kernel_without_bootstrap_providers_denies_tool_requests() {
2606        let mut kernel = Gate4AgentKernel::default();
2607        let generation = start_running(&mut kernel);
2608        let step = kernel.step_control_plane(
2609            [BackendIngress::ToolRequest(provider_bound_tool_request(
2610                None,
2611                tool_request(1, generation),
2612            ))],
2613            [],
2614        );
2615        let BackendIngressOutcome::ToolRequest(outcome) = &step.ingress_outcomes[0] else {
2616            panic!("expected tool request outcome");
2617        };
2618        assert_eq!(
2619            outcome.result,
2620            Ok(PolicyDecision::Deny(PolicyDenial::UnknownProvider))
2621        );
2622        assert!(step.tool_effects.is_empty());
2623        assert!(step.backend_snapshot.tools.providers.is_empty());
2624    }
2625
2626    #[test]
2627    fn reconciliation_failure_never_releases_retained_tool_effects() {
2628        let mut kernel =
2629            Gate4AgentKernel::with_tool_providers(builtin_registry().clone(), [tool_provider()])
2630                .unwrap();
2631        attach_tool_provider(&mut kernel, 1);
2632        start_running(&mut kernel);
2633        let divergent_generation = SessionGeneration(99);
2634        kernel
2635            .tool_engine
2636            .set_generation(instance(), divergent_generation)
2637            .unwrap();
2638        kernel
2639            .tool_engine
2640            .set_instance_state(instance(), divergent_generation, ToolInstanceState::Active)
2641            .unwrap();
2642        kernel
2643            .tool_engine
2644            .apply_authority(tool_grant(divergent_generation, 1))
2645            .unwrap();
2646        assert_eq!(
2647            kernel
2648                .tool_engine
2649                .request(tool_request(99, divergent_generation))
2650                .unwrap(),
2651            PolicyDecision::Allow
2652        );
2653
2654        let first = kernel.step_control_plane([], []);
2655        assert!(matches!(
2656            first.integration_errors[0],
2657            KernelIntegrationError::ToolInstanceSync {
2658                source: ToolEngineError::GenerationRegressed { .. },
2659                ..
2660            }
2661        ));
2662        assert!(first.tool_effects.is_empty());
2663
2664        let second = kernel.step_control_plane([], []);
2665        assert!(matches!(
2666            second.integration_errors[0],
2667            KernelIntegrationError::ToolInstanceSync {
2668                source: ToolEngineError::GenerationRegressed { .. },
2669                ..
2670            }
2671        ));
2672        assert!(second.tool_effects.is_empty());
2673    }
2674
2675    #[test]
2676    fn known_unbound_provider_is_rejected_before_canonical_acceptance() {
2677        let mut kernel =
2678            Gate4AgentKernel::with_tool_providers(builtin_registry().clone(), [tool_provider()])
2679                .unwrap();
2680        let generation = start_running(&mut kernel);
2681        let granted = kernel.step_control_plane(
2682            [BackendIngress::ToolAuthority(tool_grant(generation, 1))],
2683            [],
2684        );
2685        assert!(granted.integration_errors.is_empty());
2686
2687        let step = kernel.step_control_plane(
2688            [BackendIngress::ToolRequest(provider_bound_tool_request(
2689                None,
2690                tool_request(1, generation),
2691            ))],
2692            [],
2693        );
2694        let BackendIngressOutcome::ToolRequest(outcome) = &step.ingress_outcomes[0] else {
2695            panic!("expected tool request outcome");
2696        };
2697        assert_eq!(outcome.accepted_sequence, None);
2698        assert!(matches!(
2699            &outcome.result,
2700            Err(KernelToolError::ProviderUnavailable { provider_id })
2701                if provider_id == &tool_provider_id()
2702        ));
2703        assert!(step.backend_snapshot.tools.requests.is_empty());
2704        assert!(step.tool_effects.is_empty());
2705        assert!(step.tool_completions.completions.is_empty());
2706    }
2707
2708    #[test]
2709    fn attached_provider_round_trip_binds_effect_and_observation() {
2710        let mut kernel =
2711            Gate4AgentKernel::with_tool_providers(builtin_registry().clone(), [tool_provider()])
2712                .unwrap();
2713        let binding_id = attach_tool_provider(&mut kernel, 1);
2714        let generation = start_running(&mut kernel);
2715        kernel.step_control_plane(
2716            [BackendIngress::ToolAuthority(tool_grant(generation, 1))],
2717            [],
2718        );
2719
2720        let requested = kernel.step_control_plane(
2721            [BackendIngress::ToolRequest(provider_bound_tool_request(
2722                Some(binding_id),
2723                tool_request(1, generation),
2724            ))],
2725            [],
2726        );
2727        assert_eq!(requested.tool_effects.len(), 1);
2728        let effect = requested.tool_effects[0].clone();
2729        assert_eq!(effect.binding_id, binding_id);
2730        effect.validate().unwrap();
2731
2732        let completed = kernel.step_control_plane(
2733            [provider_runtime(
2734                2,
2735                ProviderRuntimeCommand::Observe {
2736                    binding_id,
2737                    observation: successful_observation(&effect),
2738                },
2739            )],
2740            [],
2741        );
2742        assert!(matches!(
2743            &completed.ingress_outcomes[0],
2744            BackendIngressOutcome::ToolProvider(ProviderRuntimeCommandOutcome {
2745                sequence: 2,
2746                binding_id: observed_binding,
2747                provider_id,
2748                result: Ok(ProviderRuntimeTransition::ObservationApplied { operation_id, request_key }),
2749            }) if *observed_binding == binding_id
2750                && provider_id == &tool_provider_id()
2751                && *operation_id == effect.effect.operation_id
2752                && request_key == &effect.effect.request_key
2753        ));
2754        assert!(matches!(
2755            completed.tool_completions.completions[0].outcome,
2756            CapabilityTerminalOutcome::Succeeded { .. }
2757        ));
2758        assert_eq!(completed.backend_snapshot.provider_runtime.last_sequence, 2);
2759        assert_eq!(
2760            completed.backend_snapshot.provider_runtime.bindings,
2761            vec![ProviderRuntimeBindingSnapshot {
2762                binding_id,
2763                provider_id: tool_provider_id(),
2764            }]
2765        );
2766    }
2767
2768    #[test]
2769    fn provider_observation_cannot_cross_another_active_binding() {
2770        let mut kernel = Gate4AgentKernel::with_tool_providers(
2771            builtin_registry().clone(),
2772            [tool_provider(), other_tool_provider()],
2773        )
2774        .unwrap();
2775        let first_binding = attach_tool_provider(&mut kernel, 1);
2776        let second_binding = ProviderBindingId(2);
2777        let second_attach = kernel.step_control_plane(
2778            [provider_runtime(
2779                2,
2780                ProviderRuntimeCommand::Attach {
2781                    binding_id: second_binding,
2782                    provider_id: other_tool_provider_id(),
2783                },
2784            )],
2785            [],
2786        );
2787        assert!(matches!(
2788            &second_attach.ingress_outcomes[0],
2789            BackendIngressOutcome::ToolProvider(ProviderRuntimeCommandOutcome {
2790                result: Ok(ProviderRuntimeTransition::Attached),
2791                ..
2792            })
2793        ));
2794
2795        let generation = start_running(&mut kernel);
2796        kernel.step_control_plane(
2797            [BackendIngress::ToolAuthority(tool_grant(generation, 1))],
2798            [],
2799        );
2800        let requested = kernel.step_control_plane(
2801            [BackendIngress::ToolRequest(provider_bound_tool_request(
2802                Some(first_binding),
2803                tool_request(1, generation),
2804            ))],
2805            [],
2806        );
2807        let effect = requested.tool_effects[0].clone();
2808
2809        let rejected = kernel.step_control_plane(
2810            [provider_runtime(
2811                3,
2812                ProviderRuntimeCommand::Observe {
2813                    binding_id: second_binding,
2814                    observation: successful_observation(&effect),
2815                },
2816            )],
2817            [],
2818        );
2819        assert!(matches!(
2820            &rejected.ingress_outcomes[0],
2821            BackendIngressOutcome::ToolProvider(ProviderRuntimeCommandOutcome {
2822                result: Err(KernelProviderError::BindingMismatch {
2823                    current,
2824                    requested,
2825                    ..
2826                }),
2827                ..
2828            }) if *current == first_binding && *requested == second_binding
2829        ));
2830        assert!(rejected.tool_completions.completions.is_empty());
2831
2832        let accepted = kernel.step_control_plane(
2833            [provider_runtime(
2834                4,
2835                ProviderRuntimeCommand::Observe {
2836                    binding_id: first_binding,
2837                    observation: successful_observation(&effect),
2838                },
2839            )],
2840            [],
2841        );
2842        assert!(matches!(
2843            &accepted.ingress_outcomes[0],
2844            BackendIngressOutcome::ToolProvider(ProviderRuntimeCommandOutcome {
2845                result: Ok(ProviderRuntimeTransition::ObservationApplied { .. }),
2846                ..
2847            })
2848        ));
2849    }
2850
2851    #[test]
2852    fn detach_fences_late_results_and_rebinds_without_reusing_identity() {
2853        let mut kernel =
2854            Gate4AgentKernel::with_tool_providers(builtin_registry().clone(), [tool_provider()])
2855                .unwrap();
2856        let first_binding = attach_tool_provider(&mut kernel, 1);
2857        let generation = start_running(&mut kernel);
2858        kernel.step_control_plane(
2859            [BackendIngress::ToolAuthority(tool_grant(generation, 1))],
2860            [],
2861        );
2862        let requested = kernel.step_control_plane(
2863            [BackendIngress::ToolRequest(provider_bound_tool_request(
2864                Some(first_binding),
2865                tool_request(1, generation),
2866            ))],
2867            [],
2868        );
2869        let old_effect = requested.tool_effects[0].clone();
2870
2871        let detached = kernel.step_control_plane(
2872            [provider_runtime(
2873                2,
2874                ProviderRuntimeCommand::Detach {
2875                    binding_id: first_binding,
2876                    provider_id: tool_provider_id(),
2877                },
2878            )],
2879            [],
2880        );
2881        assert!(matches!(
2882            &detached.ingress_outcomes[0],
2883            BackendIngressOutcome::ToolProvider(ProviderRuntimeCommandOutcome {
2884                result: Ok(ProviderRuntimeTransition::Detached {
2885                    closed_request_count: 1,
2886                }),
2887                ..
2888            })
2889        ));
2890        assert!(matches!(
2891            detached.tool_completions.completions[0].outcome,
2892            CapabilityTerminalOutcome::ProviderDetached { .. }
2893        ));
2894        assert!(detached
2895            .backend_snapshot
2896            .provider_runtime
2897            .bindings
2898            .is_empty());
2899        assert_eq!(detached.backend_snapshot.tools.grants.len(), 1);
2900
2901        let late = kernel.step_control_plane(
2902            [provider_runtime(
2903                3,
2904                ProviderRuntimeCommand::Observe {
2905                    binding_id: first_binding,
2906                    observation: successful_observation(&old_effect),
2907                },
2908            )],
2909            [],
2910        );
2911        assert!(matches!(
2912            &late.ingress_outcomes[0],
2913            BackendIngressOutcome::ToolProvider(ProviderRuntimeCommandOutcome {
2914                result: Err(KernelProviderError::NotAttached { .. }),
2915                ..
2916            })
2917        ));
2918        assert!(late.tool_completions.completions.is_empty());
2919
2920        let rebound = attach_tool_provider(&mut kernel, 4);
2921        assert_ne!(rebound, first_binding);
2922        let old_binding_after_rebind = kernel.step_control_plane(
2923            [provider_runtime(
2924                5,
2925                ProviderRuntimeCommand::Observe {
2926                    binding_id: first_binding,
2927                    observation: successful_observation(&old_effect),
2928                },
2929            )],
2930            [],
2931        );
2932        assert!(matches!(
2933            &old_binding_after_rebind.ingress_outcomes[0],
2934            BackendIngressOutcome::ToolProvider(ProviderRuntimeCommandOutcome {
2935                result: Err(KernelProviderError::BindingMismatch {
2936                    current,
2937                    requested,
2938                    ..
2939                }),
2940                ..
2941            }) if *current == rebound && *requested == first_binding
2942        ));
2943
2944        let next = kernel.step_control_plane(
2945            [BackendIngress::ToolRequest(provider_bound_tool_request(
2946                Some(rebound),
2947                tool_request(2, generation),
2948            ))],
2949            [],
2950        );
2951        assert_eq!(next.tool_effects[0].binding_id, rebound);
2952        let completed = kernel.step_control_plane(
2953            [provider_runtime(
2954                6,
2955                ProviderRuntimeCommand::Observe {
2956                    binding_id: rebound,
2957                    observation: successful_observation(&next.tool_effects[0]),
2958                },
2959            )],
2960            [],
2961        );
2962        assert!(matches!(
2963            completed.tool_completions.completions[0].outcome,
2964            CapabilityTerminalOutcome::Succeeded { .. }
2965        ));
2966    }
2967
2968    #[test]
2969    fn provider_request_admission_rejects_stale_missing_and_zero_bindings_without_mutation() {
2970        let mut kernel =
2971            Gate4AgentKernel::with_tool_providers(builtin_registry().clone(), [tool_provider()])
2972                .unwrap();
2973        let first_binding = attach_tool_provider(&mut kernel, 1);
2974        let detached = kernel.step_control_plane(
2975            [provider_runtime(
2976                2,
2977                ProviderRuntimeCommand::Detach {
2978                    binding_id: first_binding,
2979                    provider_id: tool_provider_id(),
2980                },
2981            )],
2982            [],
2983        );
2984        assert!(matches!(
2985            &detached.ingress_outcomes[0],
2986            BackendIngressOutcome::ToolProvider(ProviderRuntimeCommandOutcome {
2987                result: Ok(ProviderRuntimeTransition::Detached { .. }),
2988                ..
2989            })
2990        ));
2991        let rebound = attach_tool_provider(&mut kernel, 3);
2992        let generation = start_running(&mut kernel);
2993        kernel.step_control_plane(
2994            [BackendIngress::ToolAuthority(tool_grant(generation, 1))],
2995            [],
2996        );
2997
2998        let cases = [
2999            (Some(first_binding), 1_u64),
3000            (None, 2_u64),
3001            (Some(ProviderBindingId(0)), 3_u64),
3002        ];
3003        for (requested, local_id) in cases {
3004            let step = kernel.step_control_plane(
3005                [BackendIngress::ToolRequest(provider_bound_tool_request(
3006                    requested,
3007                    tool_request(local_id, generation),
3008                ))],
3009                [],
3010            );
3011            let BackendIngressOutcome::ToolRequest(outcome) = &step.ingress_outcomes[0] else {
3012                panic!("expected tool request outcome");
3013            };
3014            assert_eq!(outcome.accepted_sequence, None);
3015            if requested == Some(ProviderBindingId(0)) {
3016                assert!(matches!(
3017                    &outcome.result,
3018                    Err(KernelToolError::Validation(
3019                        ToolValidationError::ZeroIdentifier {
3020                            field: "provider binding id"
3021                        }
3022                    ))
3023                ));
3024            } else {
3025                assert!(matches!(
3026                    &outcome.result,
3027                    Err(KernelToolError::ProviderBindingMismatch {
3028                        current,
3029                        requested: rejected,
3030                        ..
3031                    }) if *current == rebound && *rejected == requested
3032                ));
3033            }
3034            assert!(step.backend_snapshot.tools.requests.is_empty());
3035            assert!(step.tool_effects.is_empty());
3036            assert!(step.tool_completions.completions.is_empty());
3037        }
3038    }
3039
3040    #[test]
3041    fn same_step_revoke_then_detach_never_releases_cancel_to_removed_binding() {
3042        let mut kernel =
3043            Gate4AgentKernel::with_tool_providers(builtin_registry().clone(), [tool_provider()])
3044                .unwrap();
3045        let binding_id = attach_tool_provider(&mut kernel, 1);
3046        let generation = start_running(&mut kernel);
3047        kernel.step_control_plane(
3048            [BackendIngress::ToolAuthority(tool_grant(generation, 1))],
3049            [],
3050        );
3051        let requested = kernel.step_control_plane(
3052            [BackendIngress::ToolRequest(provider_bound_tool_request(
3053                Some(binding_id),
3054                tool_request(1, generation),
3055            ))],
3056            [],
3057        );
3058        assert_eq!(requested.tool_effects.len(), 1);
3059
3060        let reduced = kernel.step_control_plane(
3061            [
3062                BackendIngress::ToolAuthority(ToolAuthorityEnvelope {
3063                    sequence: 2,
3064                    command: ToolAuthorityCommand::RevokeGrant {
3065                        key: tool_policy_grant(generation).key,
3066                    },
3067                }),
3068                provider_runtime(
3069                    2,
3070                    ProviderRuntimeCommand::Detach {
3071                        binding_id,
3072                        provider_id: tool_provider_id(),
3073                    },
3074                ),
3075            ],
3076            [],
3077        );
3078
3079        assert!(reduced.integration_errors.is_empty());
3080        assert!(reduced.tool_effects.is_empty());
3081        assert!(reduced
3082            .backend_snapshot
3083            .provider_runtime
3084            .bindings
3085            .is_empty());
3086        assert!(matches!(
3087            reduced.tool_completions.completions[0].outcome,
3088            CapabilityTerminalOutcome::GrantRevoked {
3089                cancellation: CancellationDisposition::CancelQueuedUnconfirmed,
3090            }
3091        ));
3092        assert!(matches!(
3093            &reduced.ingress_outcomes[1],
3094            BackendIngressOutcome::ToolProvider(ProviderRuntimeCommandOutcome {
3095                result: Ok(ProviderRuntimeTransition::Detached {
3096                    closed_request_count: 0,
3097                }),
3098                ..
3099            })
3100        ));
3101    }
3102
3103    #[test]
3104    fn provider_sequence_rejection_reuse_and_exhaustion_are_explicit() {
3105        let mut kernel =
3106            Gate4AgentKernel::with_tool_providers(builtin_registry().clone(), [tool_provider()])
3107                .unwrap();
3108        let invalid = kernel.step_control_plane(
3109            [provider_runtime(
3110                1,
3111                ProviderRuntimeCommand::Attach {
3112                    binding_id: ProviderBindingId(9),
3113                    provider_id: tool_provider_id(),
3114                },
3115            )],
3116            [],
3117        );
3118        assert!(matches!(
3119            &invalid.ingress_outcomes[0],
3120            BackendIngressOutcome::ToolProvider(ProviderRuntimeCommandOutcome {
3121                result: Err(KernelProviderError::InvalidAttachBinding { .. }),
3122                ..
3123            })
3124        ));
3125        assert_eq!(invalid.backend_snapshot.provider_runtime.last_sequence, 1);
3126
3127        let reused_sequence = kernel.step_control_plane(
3128            [provider_runtime(
3129                1,
3130                ProviderRuntimeCommand::Attach {
3131                    binding_id: ProviderBindingId(1),
3132                    provider_id: tool_provider_id(),
3133                },
3134            )],
3135            [],
3136        );
3137        assert!(matches!(
3138            &reused_sequence.ingress_outcomes[0],
3139            BackendIngressOutcome::ToolProvider(ProviderRuntimeCommandOutcome {
3140                result: Err(KernelProviderError::SequenceRegressed {
3141                    current: 1,
3142                    requested: 1,
3143                }),
3144                ..
3145            })
3146        ));
3147
3148        let binding_id = attach_tool_provider(&mut kernel, 2);
3149        let duplicate_attach = kernel.step_control_plane(
3150            [provider_runtime(
3151                3,
3152                ProviderRuntimeCommand::Attach {
3153                    binding_id: ProviderBindingId(3),
3154                    provider_id: tool_provider_id(),
3155                },
3156            )],
3157            [],
3158        );
3159        assert!(matches!(
3160            &duplicate_attach.ingress_outcomes[0],
3161            BackendIngressOutcome::ToolProvider(ProviderRuntimeCommandOutcome {
3162                result: Err(KernelProviderError::AlreadyAttached {
3163                    binding_id: current,
3164                    ..
3165                }),
3166                ..
3167            }) if *current == binding_id
3168        ));
3169        assert_eq!(
3170            duplicate_attach
3171                .backend_snapshot
3172                .provider_runtime
3173                .last_sequence,
3174            3
3175        );
3176
3177        kernel.step_control_plane(
3178            [provider_runtime(
3179                4,
3180                ProviderRuntimeCommand::Detach {
3181                    binding_id,
3182                    provider_id: tool_provider_id(),
3183                },
3184            )],
3185            [],
3186        );
3187        let reused_binding = kernel.step_control_plane(
3188            [provider_runtime(
3189                5,
3190                ProviderRuntimeCommand::Attach {
3191                    binding_id,
3192                    provider_id: tool_provider_id(),
3193                },
3194            )],
3195            [],
3196        );
3197        assert!(matches!(
3198            &reused_binding.ingress_outcomes[0],
3199            BackendIngressOutcome::ToolProvider(ProviderRuntimeCommandOutcome {
3200                result: Err(KernelProviderError::InvalidAttachBinding { .. }),
3201                ..
3202            })
3203        ));
3204
3205        let exhausted_binding = attach_tool_provider(&mut kernel, 6);
3206        assert_eq!(exhausted_binding, ProviderBindingId(6));
3207        kernel.last_provider_sequence = u64::MAX - 1;
3208        kernel.provider_sequence_exhausted = false;
3209        let unknown_provider = ToolProviderId::new("kernel-unknown-provider").unwrap();
3210        let exhausted_on_rejection = kernel.step_control_plane(
3211            [provider_runtime(
3212                u64::MAX,
3213                ProviderRuntimeCommand::Attach {
3214                    binding_id: ProviderBindingId(u64::MAX),
3215                    provider_id: unknown_provider,
3216                },
3217            )],
3218            [],
3219        );
3220        assert!(matches!(
3221            &exhausted_on_rejection.ingress_outcomes[0],
3222            BackendIngressOutcome::ToolProvider(ProviderRuntimeCommandOutcome {
3223                result: Err(KernelProviderError::UnknownProvider { .. }),
3224                ..
3225            })
3226        ));
3227        assert!(
3228            exhausted_on_rejection
3229                .backend_snapshot
3230                .provider_runtime
3231                .sequence_exhausted
3232        );
3233        assert_eq!(
3234            exhausted_on_rejection
3235                .backend_snapshot
3236                .provider_runtime
3237                .last_sequence,
3238            u64::MAX
3239        );
3240        assert!(exhausted_on_rejection
3241            .backend_snapshot
3242            .provider_runtime
3243            .bindings
3244            .is_empty());
3245
3246        let terminal = kernel.step_control_plane(
3247            [provider_runtime(
3248                u64::MAX,
3249                ProviderRuntimeCommand::Attach {
3250                    binding_id: ProviderBindingId(u64::MAX),
3251                    provider_id: tool_provider_id(),
3252                },
3253            )],
3254            [],
3255        );
3256        assert!(matches!(
3257            &terminal.ingress_outcomes[0],
3258            BackendIngressOutcome::ToolProvider(ProviderRuntimeCommandOutcome {
3259                result: Err(KernelProviderError::SequenceExhausted),
3260                ..
3261            })
3262        ));
3263    }
3264
3265    #[test]
3266    fn missing_effect_binding_is_reported_and_never_released() {
3267        let mut kernel =
3268            Gate4AgentKernel::with_tool_providers(builtin_registry().clone(), [tool_provider()])
3269                .unwrap();
3270        attach_tool_provider(&mut kernel, 1);
3271        let generation = start_running(&mut kernel);
3272        kernel
3273            .tool_engine
3274            .apply_authority(tool_grant(generation, 1))
3275            .unwrap();
3276        assert_eq!(
3277            kernel.tool_engine.request(tool_request(1, generation)),
3278            Ok(PolicyDecision::Allow)
3279        );
3280        kernel.provider_bindings.clear();
3281
3282        let step = kernel.step_control_plane([], []);
3283        assert!(matches!(
3284            &step.integration_errors[0],
3285            KernelIntegrationError::ToolEffectProviderUnbound {
3286                provider_id,
3287                ..
3288            } if provider_id == &tool_provider_id()
3289        ));
3290        assert!(step.tool_effects.is_empty());
3291    }
3292
3293    #[test]
3294    fn logical_tick_exhaustion_is_correlated_and_fail_closed() {
3295        let mut kernel = Gate4AgentKernel {
3296            logical_tick: u64::MAX,
3297            ..Gate4AgentKernel::default()
3298        };
3299        let step = kernel.step([register(1, "claude")], []);
3300
3301        assert_eq!(
3302            step.integration_errors,
3303            vec![KernelIntegrationError::LogicalTickExhausted {
3304                current_tick: u64::MAX,
3305            }]
3306        );
3307        assert!(matches!(
3308            step.command_outcomes[0].result,
3309            Err(KernelCommandError::IntegrationBlocked {
3310                reason: KernelIntegrationError::LogicalTickExhausted { .. }
3311            })
3312        ));
3313        assert!(step.snapshot.sessions.is_empty());
3314        assert!(step.effects.is_empty());
3315        assert!(step.tool_effects.is_empty());
3316        assert_eq!(step.backend_snapshot.logical_tick, u64::MAX);
3317    }
3318
3319    #[test]
3320    fn terminal_control_health_blocks_both_lanes_once() {
3321        let mut control_lane_open = true;
3322        let mut tool_lane_open = true;
3323        let mut errors = Vec::new();
3324        let healthy_capacity = ControlHealth {
3325            retained_instance_identities: 4_096,
3326            ..ControlHealth::default()
3327        };
3328
3329        block_lanes_on_control_health(
3330            healthy_capacity,
3331            &mut control_lane_open,
3332            &mut tool_lane_open,
3333            &mut errors,
3334        );
3335        assert!(control_lane_open);
3336        assert!(tool_lane_open);
3337        assert!(errors.is_empty());
3338
3339        let exhausted = ControlHealth {
3340            provider_sequence_exhausted_sessions: 1,
3341            ..healthy_capacity
3342        };
3343        block_lanes_on_control_health(
3344            exhausted,
3345            &mut control_lane_open,
3346            &mut tool_lane_open,
3347            &mut errors,
3348        );
3349        block_lanes_on_control_health(
3350            exhausted,
3351            &mut control_lane_open,
3352            &mut tool_lane_open,
3353            &mut errors,
3354        );
3355
3356        assert!(!control_lane_open);
3357        assert!(!tool_lane_open);
3358        assert_eq!(
3359            errors,
3360            vec![KernelIntegrationError::ControlHealthExhausted { health: exhausted }]
3361        );
3362    }
3363
3364    #[test]
3365    fn identical_unified_batches_produce_identical_steps() {
3366        fn run() -> KernelStep {
3367            let mut kernel = Gate4AgentKernel::with_tool_providers(
3368                builtin_registry().clone(),
3369                [tool_provider()],
3370            )
3371            .unwrap();
3372            attach_tool_provider(&mut kernel, 1);
3373            let generation = start_running(&mut kernel);
3374            kernel.step_control_plane(
3375                [
3376                    BackendIngress::ToolAuthority(tool_grant(generation, 1)),
3377                    BackendIngress::ToolRequest(provider_bound_tool_request(
3378                        Some(ProviderBindingId(1)),
3379                        tool_request(1, generation),
3380                    )),
3381                ],
3382                [],
3383            )
3384        }
3385
3386        assert_eq!(run(), run());
3387    }
3388}