1use 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#[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 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
270pub 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 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 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 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 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 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 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 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 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 #[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 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 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 #[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 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 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 ®ressed.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}