1mod efficiency;
4mod provider_supervisor;
5
6pub use efficiency::ShellEfficiencyFacts;
7pub use gate4agent::pty::ForegroundProbeTiming;
12pub use provider_supervisor::{
13 NativeProviderExecutor, NativeProviderExit, NativeProviderOperation,
14 NativeProviderOperationError, NativeProviderResultPoll, PhysicalExitAck,
15 ProviderOperationKey, ProviderOperationSnapshot, ProviderSupervisor,
16 ProviderSupervisorBuildError, ProviderSupervisorFault, ProviderSupervisorFaultKind,
17 ProviderStopCause, ProviderSupervisorSnapshot, ProviderSupervisorState,
18 ProviderSupervisorTick, DEFAULT_PROVIDER_STOP_GRACE,
19 MAX_PROVIDER_FORCE_STOP_ATTEMPTS, MAX_PROVIDER_STOP_SIGNAL_ATTEMPTS,
20 MAX_PROVIDER_SUPERVISOR_EVENTS, MAX_PROVIDER_SUPERVISOR_OPERATIONS,
21 MAX_PROVIDER_SUPERVISOR_OUTCOMES_PER_TICK, MAX_PROVIDER_SUPERVISOR_TOMBSTONES,
22 MAX_PROVIDER_SUPERVISOR_WORK_PER_TICK,
23};
24
25use gate4agent::acp::protocol::{McpServerConfig, SessionMode};
26use gate4agent::agent::{is_agent_foreground_wrapper, is_expected_agent_process, ReadinessStatus};
27use gate4agent::pty::cli::codex::strip_ansi_codes;
28use gate4agent::pty::cli::{create_pipeline, ClassificationPipeline, MessageClass, ParsedMessage};
29use gate4agent::pty::event::PtyMouseProtocolEncoding;
30use gate4agent::pty::{
31 PtyAttachment, PtyEvent, PtyEventEnvelope, PtyEventReceiver, PtyForegroundObservation,
32 PtyReplayCursor, PtySession, PtyTerminalSnapshot, RateLimitDetector, VteParser,
33};
34use gate4agent::{
35 AcpSession, AcpSessionOptions, AgentEvent, CliTool, HostDecisionAuthority, HostPolicy,
36 HostRequestDecision, HostRequestOutcome, LaunchRequest, OperatorPermissionChoice,
37 PipeProcessOptions, PipeSession, PromptFraming,
38 ReadinessIntent, ReadinessPermit, ReadinessTracker, RpcId, RuntimePlatform, SessionConfig,
39 StopReason,
40};
41use gate4agent_adapters::{
42 build_resume_plan_for_identity, builtin_adapter_registry, AdapterRuntimeRegistry,
43 CodexPtySessionIdentityExtractor, KimiPtySessionIdentityExtractor, OneShotSessionPersistence,
44};
45use gate4agent_catalog::{
46 approval_level_resolution, ApprovalLevelResolution, AgentRegistry, AgentSpec, EnvMutation,
47 McpServerSpec, ModeId,
48};
49use gate4agent_shell_one_shot::NativeOneShotSession;
50use gate4agent_types::{
51 AdapterFamily, AgentCommand, AgentId, AgentInstanceId, ApprovalLevel, CapabilityProbeFailure,
52 ContextWindowUsage as ProviderContextWindowUsage, ControlEffect,
53 ControlObservation, EffectEnvelope, ForegroundProcess, ForegroundProcessKind,
54 ForegroundRequirement, HostDecisionAuthority as ProviderHostDecisionAuthority,
55 HostRequestDecision as ProviderHostRequestDecision,
56 HostRequestOutcome as ProviderHostRequestOutcome, InputAction, ObservationEnvelope,
57 OperationId, OperatorGateInput,
58 OperatorGateKind, OperatorGateOption, OperatorGateOptionSemantics, OperatorGateState,
59 OperatorGateSubject, PipeProtocol,
60 PreparedInputKind, PromptPayload, ProviderAvailableCommand, ProviderConfigChoice,
61 ProviderConfigOption, ProviderConfigOptionKind, ProviderEvent, ProviderInteractionKind,
62 ProviderInteractionOption, ProviderInteractionResponse, ProviderInteractionTarget,
63 ProviderModeInfo, ProviderPlanPriority,
64 ProviderPlanStatus, ProviderPlanStep,
65 ProviderRateLimitKind, ProviderRuntimeCapability, ProviderRuntimePolicy,
66 ProviderSessionIdentity, ProviderSessionKey, ProviderStopReason,
67 ProviderSource, PtyScreenState, ResumeLaunchRequest, SessionGeneration, StartRequest,
68 TerminalFrame, TerminalMouseProtocolEncoding, TerminalSize, TokenUsage, TransportKind,
69 OPERATOR_GATE_OPTIONS_MAX, WORKING_DIRECTORY_MAX_BYTES,
70};
71use std::collections::{BTreeMap, VecDeque};
72use std::ffi::{OsStr, OsString};
73use std::path::PathBuf;
74use std::sync::Mutex;
75use std::time::{Duration, Instant};
76use tokio::sync::broadcast;
77use uuid::Uuid;
78
79const INSTANCE_LAUNCH_ARGS_MAX: usize = 128;
80const INSTANCE_LAUNCH_ARG_MAX_BYTES: usize = 65_536;
81const INSTANCE_LAUNCH_ARGS_TOTAL_MAX_BYTES: usize = 262_144;
82const RESERVED_CLAUDE_LAUNCH_FLAGS: &[&str] = &[
83 "--continue",
84 "--print",
85 "--prompt",
86 "--prompt-interactive",
87 "--resume",
88 "--session-id",
89 "-c",
90 "-p",
91 "-r",
92];
93
94fn mcp_server_acp_entry(spec: &McpServerSpec) -> McpServerConfig {
111 McpServerConfig::stdio(
112 spec.name().to_owned(),
113 spec.program().to_string_lossy().into_owned(),
114 spec.args()
115 .iter()
116 .map(|arg| arg.to_string_lossy().into_owned())
117 .collect::<Vec<_>>(),
118 spec.env()
119 .iter()
120 .map(|(key, value)| (key.to_string_lossy().into_owned(), value.to_string_lossy().into_owned()))
121 .collect::<Vec<_>>(),
122 )
123}
124
125#[derive(Clone, Copy, Debug, Eq, Ord, PartialEq, PartialOrd)]
126pub struct NativeSessionKey {
127 pub instance_id: AgentInstanceId,
128 pub generation: SessionGeneration,
129}
130
131struct NativeSpawnRequest {
132 agent_id: AgentId,
133 transport: TransportKind,
134 request: StartRequest,
135 runtime_policy: ProviderRuntimePolicy,
136 launch_extra_args: Vec<OsString>,
137 instance_extra_args: Vec<OsString>,
138 resumed_provider_session: Option<gate4agent_types::ProviderSessionIdentity>,
139 one_shot_session_persistence: OneShotSessionPersistence,
140}
141
142struct OwnedPtySession {
143 session: PtySession,
144 spawn_operation_id: OperationId,
145 last_terminal_sequence: u64,
146 terminal_stale_published: bool,
147 runtime_policy: ProviderRuntimePolicy,
148 provider: Option<OwnedPtyProvider>,
149 agent_id: AgentId,
154 last_screen_gate: Option<OperatorGateState>,
155 last_screen_failure: Option<&'static str>,
156 last_foreground_verdict: Option<ForegroundVerdict>,
157 last_screen_state: PtyScreenState,
158 ever_reached_ready: bool,
169 screen_had_content: bool,
174 next_foreground_probe: Option<Instant>,
179}
180
181struct OwnedPtyProvider {
182 source: ProviderSource,
183 receiver: PtyEventReceiver,
184 replay: VecDeque<PtyEventEnvelope>,
185 pending_events: VecDeque<ProviderEvent>,
186 utf8: Utf8ChunkDecoder,
187 pipeline: Mutex<ClassificationPipeline>,
188 rate_limits: RateLimitFeed,
189 kimi_identity: Option<KimiPtySessionIdentityExtractor>,
190 semantic_events: bool,
191 provider_session_started: bool,
192 next_provider_sequence: u64,
193}
194
195struct RateLimitFeed {
222 detector: RateLimitDetector,
223 vte: VteParser,
224 buffer: String,
225}
226
227impl RateLimitFeed {
228 fn new_for_tool(tool: CliTool) -> Self {
229 Self {
230 detector: RateLimitDetector::new_for_tool(tool),
231 vte: VteParser::new(),
232 buffer: String::new(),
233 }
234 }
235
236 fn detect(&mut self, raw: &str) -> Option<gate4agent::core::types::RateLimitInfo> {
246 let cleaned = self.vte.parse(raw);
247 self.buffer.push_str(&cleaned);
248 let info = self.detector.detect(&self.buffer);
249 if let Some(last_newline) = self.buffer.rfind('\n') {
250 self.buffer.drain(..=last_newline);
251 } else if self.buffer.len() > RATE_LIMIT_FEED_BUFFER_MAX_BYTES {
252 let mut cut = self.buffer.len() - RATE_LIMIT_FEED_BUFFER_MAX_BYTES;
260 while !self.buffer.is_char_boundary(cut) {
261 cut += 1;
262 }
263 self.buffer.drain(..cut);
264 }
265 info
266 }
267}
268
269const RATE_LIMIT_FEED_BUFFER_MAX_BYTES: usize = 8192;
272
273#[derive(Default)]
274struct Utf8ChunkDecoder {
275 pending: Vec<u8>,
276}
277
278impl Utf8ChunkDecoder {
279 fn push(&mut self, bytes: &[u8]) -> String {
280 self.pending.extend_from_slice(bytes);
281 let mut decoded = String::new();
282 loop {
283 match std::str::from_utf8(&self.pending) {
284 Ok(text) => {
285 decoded.push_str(text);
286 self.pending.clear();
287 break;
288 }
289 Err(error) => {
290 let valid = error.valid_up_to();
291 if valid > 0 {
292 let text = std::str::from_utf8(&self.pending[..valid])
293 .expect("validated UTF-8 prefix");
294 decoded.push_str(text);
295 self.pending.drain(..valid);
296 continue;
297 }
298 let Some(invalid) = error.error_len() else {
299 break;
300 };
301 decoded.push('\u{fffd}');
302 self.pending.drain(..invalid.min(self.pending.len()));
303 }
304 }
305 }
306 decoded
307 }
308
309 fn clear(&mut self) {
310 self.pending.clear();
311 }
312}
313
314struct OwnedProviderSession<S> {
315 source: ProviderSource,
316 session: S,
317 events: broadcast::Receiver<AgentEvent>,
318 pending_events: VecDeque<AgentEvent>,
319 pending_provider_events: VecDeque<ProviderEvent>,
328 next_provider_sequence: u64,
329 observed_exit_code: Option<i32>,
330 runtime_policy: ProviderRuntimePolicy,
331}
332
333pub struct NativeEffectShell {
336 catalog: AgentRegistry,
337 legacy_adapters: AdapterRuntimeRegistry<CliTool>,
338 pty_sessions: BTreeMap<NativeSessionKey, OwnedPtySession>,
339 pipe_sessions: BTreeMap<NativeSessionKey, OwnedProviderSession<PipeSession>>,
340 one_shot_sessions: BTreeMap<NativeSessionKey, OwnedProviderSession<NativeOneShotSession>>,
341 acp_sessions: BTreeMap<NativeSessionKey, OwnedProviderSession<AcpSession>>,
342 pending_observations: VecDeque<ObservationEnvelope>,
343 efficiency_facts: ShellEfficiencyFacts,
348}
349
350impl NativeEffectShell {
351 pub fn new(catalog: AgentRegistry) -> Self {
352 Self::new_with_runtime_adapters(catalog, builtin_legacy_adapter_runtimes())
353 }
354
355 pub fn new_with_runtime_adapters(
361 catalog: AgentRegistry,
362 legacy_adapters: AdapterRuntimeRegistry<CliTool>,
363 ) -> Self {
364 Self {
365 catalog,
366 legacy_adapters,
367 pty_sessions: BTreeMap::new(),
368 pipe_sessions: BTreeMap::new(),
369 one_shot_sessions: BTreeMap::new(),
370 acp_sessions: BTreeMap::new(),
371 pending_observations: VecDeque::new(),
372 efficiency_facts: ShellEfficiencyFacts::default(),
373 }
374 }
375
376 pub fn take_efficiency_facts(&mut self) -> ShellEfficiencyFacts {
382 self.efficiency_facts.take()
383 }
384
385 pub fn active_session_count(&self) -> usize {
386 self.pty_sessions.len()
387 + self.pipe_sessions.len()
388 + self.one_shot_sessions.len()
389 + self.acp_sessions.len()
390 }
391
392 pub fn spawn_operation_id(&self, key: NativeSessionKey) -> Option<OperationId> {
393 self.pty_sessions
394 .get(&key)
395 .map(|owned| owned.spawn_operation_id)
396 }
397
398 pub fn terminal_snapshot(&self, key: NativeSessionKey) -> Result<PtyTerminalSnapshot, String> {
399 self.pty_sessions
400 .get(&key)
401 .ok_or_else(|| missing_session_message(key))?
402 .session
403 .terminal_snapshot()
404 .map_err(|error| error.to_string())
405 }
406
407 pub async fn execute(&mut self, envelope: EffectEnvelope) -> ObservationEnvelope {
408 self.execute_with_environment(envelope, Vec::new()).await
409 }
410
411 pub async fn execute_with_pty_env(
414 &mut self,
415 envelope: EffectEnvelope,
416 pty_env: Vec<EnvMutation>,
417 ) -> ObservationEnvelope {
418 self.execute_with_environment(envelope, pty_env).await
419 }
420
421 pub async fn execute_with_environment(
425 &mut self,
426 envelope: EffectEnvelope,
427 environment: Vec<EnvMutation>,
428 ) -> ObservationEnvelope {
429 self.execute_with_launch_overlay(envelope, environment, Vec::new())
430 .await
431 }
432
433 pub async fn execute_with_launch_overlay(
436 &mut self,
437 envelope: EffectEnvelope,
438 environment: Vec<EnvMutation>,
439 extra_args: Vec<OsString>,
440 ) -> ObservationEnvelope {
441 self.execute_with_launch_overlay_and_persistence(
442 envelope,
443 environment,
444 extra_args,
445 OneShotSessionPersistence::Ephemeral,
446 )
447 .await
448 }
449
450 pub async fn execute_with_launch_overlay_and_persistence(
453 &mut self,
454 envelope: EffectEnvelope,
455 environment: Vec<EnvMutation>,
456 extra_args: Vec<OsString>,
457 one_shot_session_persistence: OneShotSessionPersistence,
458 ) -> ObservationEnvelope {
459 self.execute_with_launch_context(
460 envelope,
461 environment,
462 extra_args,
463 one_shot_session_persistence,
464 None,
465 )
466 .await
467 }
468
469 pub async fn execute_with_launch_context(
470 &mut self,
471 envelope: EffectEnvelope,
472 environment: Vec<EnvMutation>,
473 extra_args: Vec<OsString>,
474 one_shot_session_persistence: OneShotSessionPersistence,
475 mcp_server: Option<McpServerSpec>,
476 ) -> ObservationEnvelope {
477 let EffectEnvelope {
478 operation_id,
479 instance_id,
480 generation,
481 effect,
482 } = envelope;
483 let key = NativeSessionKey {
484 instance_id,
485 generation,
486 };
487
488 let observation = {
489 match effect {
490 ControlEffect::Spawn {
491 agent_id,
492 transport,
493 runtime_policy,
494 request,
495 } => {
496 self.spawn_native(
497 key,
498 operation_id,
499 NativeSpawnRequest {
500 agent_id,
501 transport,
502 request,
503 runtime_policy,
504 launch_extra_args: Vec::new(),
505 instance_extra_args: extra_args,
506 resumed_provider_session: None,
507 one_shot_session_persistence,
508 },
509 environment,
510 mcp_server,
511 )
512 .await
513 }
514 ControlEffect::SpawnResume {
515 agent_id,
516 transport,
517 provider_session,
518 runtime_policy,
519 request,
520 } => {
521 self.spawn_resume(
522 key,
523 operation_id,
524 agent_id,
525 transport,
526 provider_session,
527 runtime_policy,
528 request,
529 environment,
530 extra_args,
531 one_shot_session_persistence,
532 )
533 .await
534 }
535 ControlEffect::Stop { force } => self.stop_native(key, force).await,
536 ControlEffect::WriteInput {
537 input,
538 required_foreground,
539 } => match self.pty_sessions.get(&key) {
540 Some(owned) => match required_foreground {
541 ForegroundRequirement::Any
542 if matches!(
543 input.kind(),
544 PreparedInputKind::TerminalText
545 | PreparedInputKind::TerminalBytes
546 | PreparedInputKind::TerminalControl
547 ) =>
548 {
549 match owned.session.send_terminal_input(input).await {
550 Ok(()) => ControlObservation::InputCompleted,
551 Err(error) => ControlObservation::InputFailed {
552 message: error.to_string(),
553 },
554 }
555 }
556 ForegroundRequirement::Shell
557 if input.kind() == PreparedInputKind::ShellCommand =>
558 {
559 if !owned.runtime_policy.semantic_readiness {
560 ControlObservation::InputFailed {
561 message: "semantic shell input is not admitted by the provider runtime policy"
562 .to_owned(),
563 }
564 } else {
565 match owned.session.send_shell_input(input).await {
566 Ok(()) => ControlObservation::InputCompleted,
567 Err(error) => ControlObservation::InputFailed {
568 message: error.to_string(),
569 },
570 }
571 }
572 }
573 ForegroundRequirement::Agent { agent_id }
574 if matches!(
575 input.kind(),
576 PreparedInputKind::InsertDraft
577 | PreparedInputKind::SubmitPrompt
578 | PreparedInputKind::AgentCommand
579 ) && &agent_id == owned.session.agent_id() =>
580 {
581 if !owned.runtime_policy.semantic_readiness
582 || !owned.runtime_policy.structured_prompt
583 {
584 return completion_observation(
585 operation_id,
586 instance_id,
587 generation,
588 ControlObservation::InputFailed {
589 message: "semantic input is not admitted by the provider runtime policy"
590 .to_owned(),
591 },
592 );
593 }
594 let intent = match input.kind() {
595 PreparedInputKind::InsertDraft
596 | PreparedInputKind::AgentCommand => ReadinessIntent::DraftPaste,
597 PreparedInputKind::SubmitPrompt => ReadinessIntent::FollowupPrompt,
598 PreparedInputKind::ShellCommand
599 | PreparedInputKind::TerminalText
600 | PreparedInputKind::TerminalBytes
601 | PreparedInputKind::TerminalControl => unreachable!(),
602 };
603 let Some(spec) = self.catalog.get(&agent_id).cloned() else {
604 return completion_observation(
605 operation_id,
606 instance_id,
607 generation,
608 ControlObservation::InputFailed {
609 message: "session agent disappeared from native catalog"
610 .to_owned(),
611 },
612 );
613 };
614 match wait_for_readiness(&owned.session, &spec, intent, false).await {
615 Ok(permit) => {
616 let result = if input.kind() == PreparedInputKind::AgentCommand
617 {
618 owned.session.send_agent_command_input(input, permit).await
619 } else {
620 owned.session.send_prepared_input(input, permit).await
621 };
622 match result {
623 Ok(()) => ControlObservation::InputCompleted,
624 Err(error) => ControlObservation::InputFailed {
625 message: error.to_string(),
626 },
627 }
628 }
629 Err(message) => ControlObservation::InputFailed { message },
630 }
631 }
632 required_foreground => ControlObservation::InputFailed {
633 message: format!(
634 "prepared input kind {:?} does not satisfy route {:?}",
635 input.kind(),
636 required_foreground
637 ),
638 },
639 },
640 None => ControlObservation::InputFailed {
641 message: "typed PTY input requires a PTY session".to_owned(),
642 },
643 },
644 ControlEffect::SubmitPrompt { prompt } => match self.acp_sessions.get(&key) {
645 Some(owned) if !owned.runtime_policy.structured_prompt => {
646 ControlObservation::InputFailed {
647 message: "structured prompt is not admitted by the provider runtime policy"
648 .to_owned(),
649 }
650 }
651 Some(owned) => match owned.session.start_prompt(&prompt).await {
652 Ok(()) => ControlObservation::InputCompleted,
653 Err(error) => ControlObservation::InputFailed {
654 message: error.to_string(),
655 },
656 },
657 None => ControlObservation::InputFailed {
658 message: "semantic follow-up prompts require an ACP session".to_owned(),
659 },
660 },
661 ControlEffect::Interrupt => match self.acp_sessions.get(&key) {
662 Some(owned) => match owned.session.cancel().await {
663 Ok(()) => ControlObservation::InputCompleted,
664 Err(error) => ControlObservation::InputFailed {
665 message: error.to_string(),
666 },
667 },
668 None => ControlObservation::InputFailed {
669 message: "semantic interrupt requires an ACP session".to_owned(),
670 },
671 },
672 ControlEffect::ResolveInteraction { target, response } => {
673 resolve_acp_interaction_observation(
674 self.acp_sessions.get(&key).map(|owned| &owned.session),
675 target,
676 response,
677 )
678 }
679 ControlEffect::SetSessionMode { mode_id } => {
680 set_acp_session_mode_observation(
681 self.acp_sessions.get(&key).map(|owned| &owned.session),
682 mode_id,
683 )
684 .await
685 }
686 ControlEffect::SetSessionConfigOption {
687 option_id,
688 value_json,
689 } => {
690 set_acp_session_config_option_observation(
691 self.acp_sessions.get(&key).map(|owned| &owned.session),
692 option_id,
693 value_json,
694 )
695 .await
696 }
697 ControlEffect::SetSessionModel { model_id } => {
698 set_acp_session_model_observation(
699 self.acp_sessions.get(&key).map(|owned| &owned.session),
700 model_id,
701 )
702 }
703 ControlEffect::Resize { size } if !size.is_valid() => {
704 ControlObservation::ResizeFailed {
705 message: "terminal size is outside the supported range".to_owned(),
706 }
707 }
708 ControlEffect::Resize { size } => match self.pty_sessions.get(&key) {
709 Some(owned) => match owned.session.resize(size.rows, size.columns).await {
710 Ok(()) => ControlObservation::ResizeCompleted { size },
711 Err(error) => ControlObservation::ResizeFailed {
712 message: error.to_string(),
713 },
714 },
715 None => ControlObservation::ResizeFailed {
716 message: "terminal resize requires a PTY session".to_owned(),
717 },
718 },
719 ControlEffect::ObserveForeground => match self.pty_sessions.get(&key) {
720 Some(owned) => match owned.session.observe_foreground().await {
721 Ok(observation) => ControlObservation::ForegroundObserved {
722 process: canonical_foreground(
723 owned.session.agent_id().clone(),
724 &observation,
725 ),
726 },
727 Err(error) => ControlObservation::ForegroundFailed {
728 message: error.to_string(),
729 },
730 },
731 None => ControlObservation::ForegroundFailed {
732 message: "foreground observation requires a PTY session".to_owned(),
733 },
734 },
735 ControlEffect::ProbeCapabilities { .. } => {
736 ControlObservation::CapabilityProbeFailed {
737 failure: CapabilityProbeFailure::ExecutorUnavailable,
738 }
739 }
740 ControlEffect::DiscoverHistory { .. } | ControlEffect::LoadHistory { .. } => {
741 ControlObservation::HistoryFailed {
742 message: "history effects require the dedicated native history authority"
743 .to_owned(),
744 }
745 }
746 ControlEffect::AuthorizeResume { .. } => ControlObservation::ResumeFailed {
747 message: "resume authorization requires the dedicated native authority"
748 .to_owned(),
749 },
750 }
751 };
752
753 completion_observation(operation_id, instance_id, generation, observation)
754 }
755
756 async fn spawn_native(
757 &mut self,
758 key: NativeSessionKey,
759 operation_id: OperationId,
760 spawn: NativeSpawnRequest,
761 pty_env: Vec<EnvMutation>,
762 mcp_server: Option<McpServerSpec>,
763 ) -> ControlObservation {
764 let NativeSpawnRequest {
765 agent_id,
766 transport,
767 request,
768 runtime_policy,
769 mut launch_extra_args,
770 instance_extra_args,
771 resumed_provider_session,
772 one_shot_session_persistence,
773 } = spawn;
774 if let Err(message) =
775 validate_instance_launch_arguments(&agent_id, transport, &instance_extra_args)
776 {
777 return ControlObservation::SpawnFailed { message };
778 }
779 if let Err(message) = validate_spawn_runtime_policy(
780 runtime_policy,
781 transport,
782 request.initial_prompt.is_some(),
783 resumed_provider_session.is_some(),
784 ) {
785 return ControlObservation::SpawnFailed { message };
786 }
787 if self.session_exists(key) {
788 return ControlObservation::SpawnFailed {
789 message: format!(
790 "native session {:?}/{:?} already exists",
791 key.instance_id, key.generation
792 ),
793 };
794 }
795 if !request.terminal_size.is_valid() {
796 return ControlObservation::SpawnFailed {
797 message: "terminal size is outside the supported range".to_owned(),
798 };
799 }
800 if request.working_directory.is_empty()
801 || request.working_directory.len() > WORKING_DIRECTORY_MAX_BYTES
802 || request.working_directory.contains('\0')
803 {
804 return ControlObservation::SpawnFailed {
805 message: "working directory is invalid".to_owned(),
806 };
807 }
808 let Some(spec) = self.catalog.get(&agent_id).cloned() else {
809 return ControlObservation::SpawnFailed {
810 message: format!("agent '{agent_id}' is absent from native catalog"),
811 };
812 };
813 let working_dir = PathBuf::from(&request.working_directory);
814
815 match transport {
816 TransportKind::Pty if !spec.capabilities.transports.pty => {
817 ControlObservation::SpawnFailed {
818 message: format!("agent '{agent_id}' does not support PTY transport"),
819 }
820 }
821 TransportKind::Pty => {
822 let pty_env = with_pty_terminal_capability_defaults(pty_env);
827 let fresh_provider_session = prepare_fresh_pty_provider_session(
828 spec.capabilities.transports.pty_adapter.as_ref(),
829 resumed_provider_session.is_some(),
830 runtime_policy.provider_session_identity,
831 &mut launch_extra_args,
832 );
833 launch_extra_args.extend(instance_extra_args);
834 let mut authoritative_provider_session = resumed_provider_session
835 .clone()
836 .or(fresh_provider_session);
837 let probe_kimi_identity = should_probe_pty_identity(
838 runtime_policy,
839 spec.capabilities.transports.pty_adapter.as_ref(),
840 authoritative_provider_session.is_some(),
841 "kimi",
842 );
843 let probe_codex_identity = should_probe_pty_identity(
844 runtime_policy,
845 spec.capabilities.transports.pty_adapter.as_ref(),
846 authoritative_provider_session.is_some(),
847 "codex",
848 );
849 tracing::info!(
862 agent_id = %agent_id,
863 instance_id = ?key.instance_id,
864 generation = ?key.generation,
865 terminal_size = ?request.terminal_size,
866 program = %spec.launch.program,
867 fixed_args = ?spec.launch.fixed_args,
868 extra_args = ?redact_provider_arguments(&launch_extra_args),
869 has_initial_prompt = request.initial_prompt.is_some(),
870 "spawning provider process over a PTY",
871 );
872 match PtySession::spawn_agent_with_size(
873 &spec,
874 LaunchRequest {
875 working_dir,
876 env: pty_env,
877 platform: RuntimePlatform::current(),
878 prompt: request.initial_prompt,
879 session_options: request.session_options,
880 extra_args: launch_extra_args,
881 },
882 request.terminal_size.rows,
883 request.terminal_size.columns,
884 )
885 .await
886 {
887 Ok(mut session) => {
888 if probe_kimi_identity {
889 match probe_fresh_kimi_session_identity(&session, &spec).await {
890 Ok(Some(identity)) => authoritative_provider_session = Some(identity),
891 Ok(None) => {}
892 Err(error) => {
893 let message = match session.shutdown().await {
894 Ok(_) => error,
895 Err(shutdown_error) => {
896 format!("{error}; PTY cleanup failed: {shutdown_error}")
897 }
898 };
899 return ControlObservation::SpawnFailed { message };
900 }
901 }
902 }
903 if probe_codex_identity {
904 match probe_fresh_codex_session_identity(&session, &spec).await {
905 Ok(Some(identity)) => authoritative_provider_session = Some(identity),
906 Ok(None) => {}
907 Err(error) => {
908 let message = match session.shutdown().await {
909 Ok(_) => error,
910 Err(shutdown_error) => {
911 format!("{error}; PTY cleanup failed: {shutdown_error}")
912 }
913 };
914 return ControlObservation::SpawnFailed { message };
915 }
916 }
917 }
918 if let Err(error) = deliver_pending_initial_prompt(&mut session, &spec).await {
919 let message = match session.shutdown().await {
920 Ok(_) => error,
921 Err(shutdown_error) => {
922 format!("{error}; PTY cleanup failed: {shutdown_error}")
923 }
924 };
925 return ControlObservation::SpawnFailed { message };
926 }
927 let process_id = session.root_pid();
928 let provider = match spec.capabilities.transports.pty_adapter.as_ref() {
929 _ if !should_attach_pty_provider_stream(runtime_policy) => None,
930 Some(adapter) => {
931 let tool = match self
932 .legacy_adapters
933 .resolve(AdapterFamily::PtySemantic, adapter)
934 {
935 Ok(tool) => *tool,
936 Err(error) => {
937 let _ = session.shutdown().await;
938 return ControlObservation::SpawnFailed {
939 message: error.to_string(),
940 };
941 }
942 };
943 match session.attach_events(session.beginning_cursor()) {
944 Ok(attachment) => {
945 let mut pending_events = VecDeque::new();
946 let mut provider_session_started = false;
947 if let Some(identity) = authoritative_provider_session {
948 pending_events.push_back(ProviderEvent::SessionStarted {
949 session_id: identity.id.clone(),
950 model: String::new(),
951 tools: Vec::new(),
952 });
953 pending_events.push_back(
954 ProviderEvent::SessionIdentityObserved { identity },
955 );
956 provider_session_started = true;
957 }
958 let is_kimi = tool == CliTool::KimiCode;
959 Some(OwnedPtyProvider {
960 source: ProviderSource {
961 family: AdapterFamily::PtySemantic,
962 binding: adapter.clone(),
963 },
964 receiver: attachment.receiver,
965 replay: attachment.replay.into(),
966 pending_events,
967 utf8: Utf8ChunkDecoder::default(),
968 pipeline: Mutex::new(create_pipeline(tool)),
969 rate_limits: RateLimitFeed::new_for_tool(tool),
970 kimi_identity: (runtime_policy.provider_session_identity
971 && is_kimi
972 && !provider_session_started)
973 .then(KimiPtySessionIdentityExtractor::default),
974 semantic_events: runtime_policy.semantic_readiness,
975 provider_session_started,
976 next_provider_sequence: 1,
977 })
978 }
979 Err(error) => {
980 let _ = session.shutdown().await;
981 return ControlObservation::SpawnFailed {
982 message: error.to_string(),
983 };
984 }
985 }
986 }
987 None => None,
988 };
989 self.pty_sessions.insert(
990 key,
991 OwnedPtySession {
992 session,
993 spawn_operation_id: operation_id,
994 last_terminal_sequence: 0,
995 terminal_stale_published: false,
996 runtime_policy,
997 provider,
998 agent_id,
999 last_screen_gate: None,
1000 last_screen_failure: None,
1001 last_foreground_verdict: None,
1002 last_screen_state: PtyScreenState::default(),
1003 ever_reached_ready: false,
1004 screen_had_content: false,
1005 next_foreground_probe: Some(Instant::now()),
1010 },
1011 );
1012 ControlObservation::Spawned { process_id }
1013 }
1014 Err(error) => {
1015 let message = error.to_string();
1016 tracing::warn!(
1017 agent_id = %agent_id,
1018 instance_id = ?key.instance_id,
1019 generation = ?key.generation,
1020 program = %spec.launch.program,
1021 fixed_args = ?spec.launch.fixed_args,
1022 cause = %message,
1023 "provider process failed to start",
1024 );
1025 ControlObservation::SpawnFailed { message }
1026 }
1027 }
1028 }
1029 TransportKind::Pipe => {
1030 let Some(pipe_spec) = spec.capabilities.transports.pipe.as_ref() else {
1031 return ControlObservation::SpawnFailed {
1032 message: format!("agent '{agent_id}' does not support Pipe transport"),
1033 };
1034 };
1035 let prompt = request.initial_prompt.unwrap_or_default();
1036 if pipe_spec.protocol == PipeProtocol::OneShotText {
1037 let Some(binding) = spec.capabilities.adapters.one_shot.as_ref() else {
1038 return ControlObservation::SpawnFailed {
1039 message: format!(
1040 "agent '{agent_id}' does not declare OneShot capability"
1041 ),
1042 };
1043 };
1044 if binding != &pipe_spec.adapter {
1045 return ControlObservation::SpawnFailed {
1046 message: format!(
1047 "agent '{agent_id}' has mismatched OneShot transport bindings"
1048 ),
1049 };
1050 }
1051 return match NativeOneShotSession::spawn_with_environment_and_persistence(
1052 &spec,
1053 binding,
1054 &prompt,
1055 request.session_options.as_ref(),
1056 &working_dir,
1057 &pty_env,
1058 one_shot_session_persistence,
1059 )
1060 .await
1061 {
1062 Ok(session) => {
1063 let process_id = session.process_id();
1064 let events = session.subscribe();
1065 self.one_shot_sessions.insert(
1066 key,
1067 OwnedProviderSession {
1068 source: ProviderSource {
1069 family: AdapterFamily::OneShot,
1070 binding: binding.clone(),
1071 },
1072 session,
1073 events,
1074 pending_events: VecDeque::new(),
1075 pending_provider_events: VecDeque::new(),
1076 next_provider_sequence: 1,
1077 observed_exit_code: None,
1078 runtime_policy,
1079 },
1080 );
1081 ControlObservation::Spawned { process_id }
1082 }
1083 Err(error) => ControlObservation::SpawnFailed {
1084 message: error.to_string(),
1085 },
1086 };
1087 }
1088 let tool = match self
1089 .legacy_adapters
1090 .resolve(AdapterFamily::Pipe, &pipe_spec.adapter)
1091 {
1092 Ok(tool) => *tool,
1093 Err(error) => {
1094 return ControlObservation::SpawnFailed {
1095 message: error.to_string(),
1096 }
1097 }
1098 };
1099 let source = ProviderSource {
1100 family: AdapterFamily::Pipe,
1101 binding: pipe_spec.adapter.clone(),
1102 };
1103 let config = SessionConfig {
1104 tool,
1105 working_dir,
1106 env_vars: Vec::new(),
1107 name: None,
1108 };
1109 if resumed_provider_session.is_some() && pipe_spec.launch_override.is_some() {
1110 return ControlObservation::SpawnFailed {
1111 message: format!(
1112 "agent '{agent_id}' cannot resume through a catalog launch override"
1113 ),
1114 };
1115 }
1116 let mut options = PipeProcessOptions::default();
1117 if let Some(identity) = resumed_provider_session.as_ref() {
1118 options.claude.resume_session_id = Some(identity.id.clone());
1119 }
1120 let spawned = match pipe_spec.launch_override.as_ref() {
1121 Some(launch) => {
1122 PipeSession::spawn_with_launch(
1123 config,
1124 &prompt,
1125 launch,
1126 pipe_spec.prompt_delivery,
1127 )
1128 .await
1129 }
1130 None => {
1131 PipeSession::spawn(config, &prompt, options).await
1132 }
1133 };
1134 match spawned {
1135 Ok(session) => {
1136 let process_id = session.process_id();
1137 let events = session.subscribe();
1138 let pending_events = if pipe_spec.protocol == PipeProtocol::SemanticNdjson {
1139 VecDeque::from([AgentEvent::SessionStart {
1140 session_id: session.session_id().to_owned(),
1141 model: String::new(),
1142 tools: Vec::new(),
1143 }])
1144 } else {
1145 VecDeque::new()
1146 };
1147 self.pipe_sessions.insert(
1148 key,
1149 OwnedProviderSession {
1150 source,
1151 session,
1152 events,
1153 pending_events,
1154 pending_provider_events: VecDeque::new(),
1155 next_provider_sequence: 1,
1156 observed_exit_code: None,
1157 runtime_policy,
1158 },
1159 );
1160 ControlObservation::Spawned { process_id }
1161 }
1162 Err(error) => ControlObservation::SpawnFailed {
1163 message: error.to_string(),
1164 },
1165 }
1166 }
1167 TransportKind::Acp => {
1168 let Some(acp_spec) = spec.capabilities.transports.acp else {
1169 return ControlObservation::SpawnFailed {
1170 message: format!("agent '{agent_id}' does not support ACP transport"),
1171 };
1172 };
1173 let tool = match self
1174 .legacy_adapters
1175 .resolve(AdapterFamily::Acp, &acp_spec.adapter)
1176 {
1177 Ok(tool) => *tool,
1178 Err(error) => {
1179 return ControlObservation::SpawnFailed {
1180 message: error.to_string(),
1181 }
1182 }
1183 };
1184 let source = ProviderSource {
1185 family: AdapterFamily::Acp,
1186 binding: acp_spec.adapter.clone(),
1187 };
1188 let acp_mode_id = match required_acp_mode(&agent_id, request.approval_level) {
1217 Ok(mode_id) => mode_id,
1218 Err(message) => return ControlObservation::SpawnFailed { message },
1219 };
1220 let mcp_servers: Vec<McpServerConfig> =
1228 mcp_server.as_ref().map(mcp_server_acp_entry).into_iter().collect();
1229 if !mcp_servers.is_empty() {
1230 tracing::info!(
1238 agent_id = %agent_id,
1239 instance_id = ?key.instance_id,
1240 mcp_server_names = ?mcp_servers.iter().map(McpServerConfig::name).collect::<Vec<_>>(),
1241 "acp session/new carries {} mcp server(s)",
1242 mcp_servers.len(),
1243 );
1244 }
1245 let mut approval_level_args =
1249 acp_approval_level_args(&agent_id, request.approval_level);
1250 for argument in &instance_extra_args {
1251 approval_level_args.push(argument.to_string_lossy().into_owned());
1252 }
1253 let acp_options = AcpSessionOptions {
1254 host_policy: host_policy_for_approval_level(request.approval_level),
1255 approval_level_args,
1270 defer_permission_requests: defers_permission_requests(
1271 &agent_id,
1272 request.approval_level,
1273 ),
1274 mcp_servers,
1275 ..AcpSessionOptions::default()
1276 };
1277 let spawned = match acp_spec.launch_override.as_ref() {
1278 Some(launch) => {
1279 AcpSession::spawn_with_launch(tool, &working_dir, acp_options, launch)
1280 .await
1281 }
1282 None => AcpSession::spawn(tool, &working_dir, acp_options).await,
1283 };
1284 match spawned {
1285 Ok(session) => {
1286 if let Some(mode_id) = acp_mode_id.as_ref() {
1287 if let Err(message) =
1288 apply_acp_approval_mode(&session, mode_id, request.approval_level)
1289 .await
1290 {
1291 let _ = session.kill().await;
1292 return ControlObservation::SpawnFailed { message };
1293 }
1294 }
1295 let process_id = session.process_id();
1296 let events = session.subscribe();
1297 let session_id = session
1298 .acp_session_id()
1299 .await
1300 .unwrap_or_else(|| session.session_id().to_owned());
1301 let identity_observed = if runtime_policy.provider_session_identity {
1331 match session.provider_reported_session_id().await {
1332 Some(id) => Some(ProviderEvent::SessionIdentityObserved {
1333 identity: ProviderSessionIdentity {
1334 key: ProviderSessionKey::SessionId,
1335 id,
1336 transcript_path: None,
1337 },
1338 }),
1339 None => {
1340 tracing::warn!(
1341 agent_id = %agent_id,
1342 adapter = %acp_spec.adapter.id,
1343 "ACP session/new returned no sessionId; provider \
1344 session identity stays IdentityPending",
1345 );
1346 None
1347 }
1348 }
1349 } else {
1350 None
1351 };
1352 if let Some(prompt) = request.initial_prompt {
1353 if let Err(error) = session.start_prompt(&prompt).await {
1354 let _ = session.kill().await;
1355 return ControlObservation::SpawnFailed {
1356 message: error.to_string(),
1357 };
1358 }
1359 }
1360 let announced_mode = session.current_mode_id();
1362 self.acp_sessions.insert(
1363 key,
1364 OwnedProviderSession {
1365 source,
1366 session,
1367 events,
1368 pending_events: {
1369 let mut seeded = VecDeque::from([AgentEvent::SessionStart {
1370 session_id,
1371 model: String::new(),
1372 tools: Vec::new(),
1373 }]);
1374 if let Some(mode_id) = announced_mode {
1397 seeded.push_back(AgentEvent::ModeChanged { mode_id });
1398 }
1399 seeded
1400 },
1401 pending_provider_events: VecDeque::from_iter(identity_observed),
1402 next_provider_sequence: 1,
1403 observed_exit_code: None,
1404 runtime_policy,
1405 },
1406 );
1407 ControlObservation::Spawned { process_id }
1408 }
1409 Err(error) => ControlObservation::SpawnFailed {
1410 message: error.to_string(),
1411 },
1412 }
1413 }
1414 }
1415 }
1416
1417 async fn spawn_resume(
1418 &mut self,
1419 key: NativeSessionKey,
1420 operation_id: OperationId,
1421 agent_id: AgentId,
1422 transport: TransportKind,
1423 provider_session: gate4agent_types::ProviderSessionIdentity,
1424 runtime_policy: ProviderRuntimePolicy,
1425 request: ResumeLaunchRequest,
1426 pty_env: Vec<EnvMutation>,
1427 instance_extra_args: Vec<OsString>,
1428 one_shot_session_persistence: OneShotSessionPersistence,
1429 ) -> ControlObservation {
1430 if let Err(message) = validate_spawn_runtime_policy(
1431 runtime_policy,
1432 transport,
1433 request.initial_prompt.is_some(),
1434 true,
1435 ) {
1436 return ControlObservation::SpawnFailed { message };
1437 }
1438 if let Err(error) = request.validate() {
1439 return ControlObservation::SpawnFailed {
1440 message: error.to_string(),
1441 };
1442 }
1443 let Some(spec) = self.catalog.get(&agent_id) else {
1444 return ControlObservation::SpawnFailed {
1445 message: format!("agent '{agent_id}' is absent from native catalog"),
1446 };
1447 };
1448 let Some(binding) = spec.capabilities.adapters.resume.as_ref() else {
1449 return ControlObservation::SpawnFailed {
1450 message: format!("agent '{agent_id}' does not declare Resume capability"),
1451 };
1452 };
1453 let plan = match build_resume_plan_for_identity(&binding.id, &provider_session) {
1454 Ok(Some(plan)) => plan,
1455 Ok(None) => {
1456 return ControlObservation::SpawnFailed {
1457 message: format!("agent '{agent_id}' has no live Resume plan"),
1458 }
1459 }
1460 Err(error) => {
1461 return ControlObservation::SpawnFailed {
1462 message: error.to_string(),
1463 }
1464 }
1465 };
1466 if transport == TransportKind::Pipe {
1467 let Some(pipe) = spec.capabilities.transports.pipe.as_ref() else {
1468 return ControlObservation::SpawnFailed {
1469 message: format!("agent '{agent_id}' does not support Pipe transport"),
1470 };
1471 };
1472 if pipe.protocol != PipeProtocol::StructuredJsonl {
1473 return ControlObservation::SpawnFailed {
1474 message: format!(
1475 "agent '{agent_id}' does not expose a resumable structured Pipe contract"
1476 ),
1477 };
1478 }
1479 }
1480 let start = StartRequest {
1481 working_directory: request.working_directory,
1482 terminal_size: request.terminal_size,
1483 initial_prompt: request.initial_prompt,
1484 session_options: None,
1490 approval_level: ApprovalLevel::default(),
1491 };
1492 self.spawn_native(
1493 key,
1494 operation_id,
1495 NativeSpawnRequest {
1496 agent_id,
1497 transport,
1498 request: start,
1499 runtime_policy,
1500 launch_extra_args: if transport == TransportKind::Pty {
1501 plan.args.into_iter().map(OsString::from).collect()
1502 } else {
1503 Vec::new()
1504 },
1505 instance_extra_args,
1506 resumed_provider_session: Some(provider_session),
1507 one_shot_session_persistence,
1508 },
1509 pty_env,
1510 None,
1511 )
1512 .await
1513 }
1514
1515 async fn stop_native(&mut self, key: NativeSessionKey, force: bool) -> ControlObservation {
1516 if let Some(owned) = self.pty_sessions.remove(&key) {
1517 let shutdown = owned.session.shutdown().await;
1518 return match shutdown {
1519 Ok(outcome) => ControlObservation::StopCompleted {
1520 forced: force || outcome.termination.is_some(),
1521 exit_code: outcome.exit_code,
1522 final_terminal: Some(terminal_frame(outcome.terminal, owned.last_screen_state.clone())),
1526 },
1527 Err(error) => {
1528 eprintln!(
1529 "[gate4agent-shell-native] PTY stop failed for {key:?} (force={force}): {error}",
1530 );
1531 ControlObservation::StopFailed {
1532 message: error.to_string(),
1533 }
1534 }
1535 };
1536 }
1537 if let Some(owned) = self.pipe_sessions.remove(&key) {
1538 return match owned.session.stop(force).await {
1551 Ok(outcome) => ControlObservation::StopCompleted {
1552 forced: outcome.forced,
1553 exit_code: outcome.exit_code.or(owned.observed_exit_code),
1554 final_terminal: None,
1555 },
1556 Err(_) if owned.session.reader_finished() => ControlObservation::StopCompleted {
1557 forced: false,
1558 exit_code: owned.observed_exit_code,
1559 final_terminal: None,
1560 },
1561 Err(error) => ControlObservation::StopFailed {
1562 message: error.to_string(),
1563 },
1564 };
1565 }
1566 if let Some(mut owned) = self.one_shot_sessions.remove(&key) {
1567 return match owned.session.kill().await {
1568 Ok(()) => ControlObservation::StopCompleted {
1569 forced: true,
1570 exit_code: owned.observed_exit_code,
1571 final_terminal: None,
1572 },
1573 Err(_) if owned.session.reader_finished() => ControlObservation::StopCompleted {
1574 forced: false,
1575 exit_code: owned.observed_exit_code,
1576 final_terminal: None,
1577 },
1578 Err(error) => ControlObservation::StopFailed {
1579 message: error.to_string(),
1580 },
1581 };
1582 }
1583 if let Some(owned) = self.acp_sessions.remove(&key) {
1584 return match owned.session.stop(force).await {
1595 Ok(outcome) => ControlObservation::StopCompleted {
1596 forced: outcome.forced,
1597 exit_code: outcome.exit_code.or(owned.observed_exit_code),
1598 final_terminal: None,
1599 },
1600 Err(_) if owned.session.reader_finished() => ControlObservation::StopCompleted {
1601 forced: false,
1602 exit_code: owned.observed_exit_code,
1603 final_terminal: None,
1604 },
1605 Err(error) => ControlObservation::StopFailed {
1606 message: error.to_string(),
1607 },
1608 };
1609 }
1610 eprintln!("[gate4agent-shell-native] stop skipped: no session owns {key:?} in any transport map");
1616 ControlObservation::StopFailed {
1617 message: missing_session_message(key),
1618 }
1619 }
1620
1621 fn session_exists(&self, key: NativeSessionKey) -> bool {
1622 self.pty_sessions.contains_key(&key)
1623 || self.pipe_sessions.contains_key(&key)
1624 || self.one_shot_sessions.contains_key(&key)
1625 || self.acp_sessions.contains_key(&key)
1626 }
1627
1628 pub async fn collect_exits(&mut self) -> Vec<ObservationEnvelope> {
1631 let completed: Vec<_> = self
1632 .pty_sessions
1633 .iter()
1634 .filter_map(|(key, owned)| owned.session.reader_finished().then_some(*key))
1635 .collect();
1636 let mut observations = Vec::with_capacity(completed.len());
1637 for key in completed {
1638 let owned = self
1639 .pty_sessions
1640 .remove(&key)
1641 .expect("completed key came from the owned session map");
1642 let last_screen_state = owned.last_screen_state.clone();
1643 let (exit_code, final_terminal) = match owned.session.shutdown().await {
1644 Ok(outcome) => (
1645 outcome.exit_code,
1646 Some(terminal_frame(outcome.terminal, last_screen_state)),
1647 ),
1648 Err(_) => (None, None),
1649 };
1650 observations.push(ObservationEnvelope {
1651 operation_id: None,
1652 instance_id: key.instance_id,
1653 generation: key.generation,
1654 observation: ControlObservation::ProcessExited {
1655 exit_code,
1656 final_terminal,
1657 },
1658 });
1659 }
1660 collect_provider_exits(&mut self.pipe_sessions, &mut observations, |session| {
1661 session.reader_finished()
1662 });
1663 collect_provider_exits(&mut self.one_shot_sessions, &mut observations, |session| {
1664 session.reader_finished()
1665 });
1666 collect_provider_exits(&mut self.acp_sessions, &mut observations, |session| {
1667 session.reader_finished()
1668 });
1669 observations
1670 }
1671
1672 pub fn collect_provider_events(&mut self) -> Vec<ObservationEnvelope> {
1675 let mut observations = self.pending_observations.drain(..).collect::<Vec<_>>();
1676 for (key, owned) in &mut self.pty_sessions {
1677 if let Some(provider) = &mut owned.provider {
1678 drain_pty_provider(*key, provider, &mut observations);
1679 }
1680 }
1681 collect_provider_map(&mut self.pipe_sessions, &mut observations, |_| Vec::new());
1682 collect_provider_map(&mut self.one_shot_sessions, &mut observations, |_| Vec::new());
1683 for owned in self.acp_sessions.values() {
1695 owned.session.expire_deadlines();
1696 }
1697 collect_provider_map(&mut self.acp_sessions, &mut observations, AcpSession::available_modes);
1698 observations
1699 }
1700
1701 pub fn collect_terminal_frames(&mut self) -> Vec<ObservationEnvelope> {
1715 let mut observations = Vec::new();
1716 let pty_sessions = &mut self.pty_sessions;
1722 let efficiency_facts = &mut self.efficiency_facts;
1723 for (key, owned) in pty_sessions {
1724 if terminal_state_capture_should_skip(
1725 owned.session.terminal_sequence(),
1726 owned.last_terminal_sequence,
1727 ) {
1728 efficiency_facts.record_terminal_state_skip();
1729 continue;
1730 }
1731 let capture_start = Instant::now();
1732 let terminal_state = owned.session.terminal_state();
1733 efficiency_facts.record_terminal_state_capture(capture_start.elapsed());
1734 match terminal_state {
1735 Ok(snapshot) if snapshot.sequence > owned.last_terminal_sequence => {
1736 owned.last_terminal_sequence = snapshot.sequence;
1737 owned.terminal_stale_published = false;
1738
1739 owned.last_screen_gate = startup_operator_gate(&snapshot.contents);
1743 owned.last_screen_failure =
1754 screen_failure_for_generation(&snapshot.contents, owned.ever_reached_ready);
1755 let merged = classify_pty_screen_state(
1756 owned.last_foreground_verdict.as_ref(),
1757 owned.last_screen_gate.as_ref(),
1758 owned.last_screen_failure,
1759 );
1760 if !snapshot.contents.trim().is_empty() {
1771 owned.screen_had_content = true;
1772 }
1773 if merged == PtyScreenState::Ready && owned.screen_had_content {
1774 owned.ever_reached_ready = true;
1775 }
1776 if merged != owned.last_screen_state {
1777 if foreground_probe_rearms_immediately(&owned.last_screen_state, &merged) {
1778 owned.next_foreground_probe = Some(Instant::now());
1779 }
1780 owned.last_screen_state = merged.clone();
1781 observations.push(ObservationEnvelope {
1782 operation_id: None,
1783 instance_id: key.instance_id,
1784 generation: key.generation,
1785 observation: ControlObservation::ScreenState {
1786 state: merged.clone(),
1787 },
1788 });
1789 }
1790
1791 let frame = terminal_frame(snapshot, merged);
1792 efficiency_facts.record_terminal_frame_published(terminal_frame_byte_len(&frame));
1793 observations.push(ObservationEnvelope {
1794 operation_id: None,
1795 instance_id: key.instance_id,
1796 generation: key.generation,
1797 observation: ControlObservation::TerminalFrame { frame },
1798 });
1799 }
1800 Ok(_) => {}
1801 Err(error) if !owned.terminal_stale_published => {
1802 owned.terminal_stale_published = true;
1803 observations.push(ObservationEnvelope {
1804 operation_id: None,
1805 instance_id: key.instance_id,
1806 generation: key.generation,
1807 observation: ControlObservation::TerminalStale {
1808 message: error.to_string(),
1809 },
1810 });
1811 }
1812 Err(_) => {}
1813 }
1814 }
1815 observations
1816 }
1817
1818 pub async fn reclassify_foreground(&mut self) -> Vec<ObservationEnvelope> {
1836 let catalog = &self.catalog;
1837 let now = Instant::now();
1838 let due: Vec<NativeSessionKey> = self
1839 .pty_sessions
1840 .iter()
1841 .filter(|(_, owned)| owned.next_foreground_probe.is_some_and(|at| now >= at))
1842 .map(|(key, _)| *key)
1843 .collect();
1844
1845 let mut observations = Vec::new();
1846 for key in due {
1847 let Some(owned) = self.pty_sessions.get_mut(&key) else {
1848 continue;
1849 };
1850 let probe_result = owned.session.observe_foreground_timed().await;
1851 match probe_result {
1852 Ok((observation, timing)) => {
1853 self.efficiency_facts.record_foreground_probe(timing);
1866 let verdict = match catalog.get(&owned.agent_id) {
1867 Some(spec) => resolve_foreground_verdict(
1868 spec,
1869 &observation,
1870 RuntimePlatform::current(),
1871 ),
1872 None => ForegroundVerdict::Foreign {
1878 process: observation.observed_process.clone(),
1879 },
1880 };
1881 owned.last_foreground_verdict = Some(verdict);
1882 let merged = classify_pty_screen_state(
1883 owned.last_foreground_verdict.as_ref(),
1884 owned.last_screen_gate.as_ref(),
1885 owned.last_screen_failure,
1886 );
1887 if merged == PtyScreenState::Ready && owned.screen_had_content {
1892 owned.ever_reached_ready = true;
1893 }
1894 if merged != owned.last_screen_state {
1895 owned.last_screen_state = merged.clone();
1896 observations.push(ObservationEnvelope {
1897 operation_id: None,
1898 instance_id: key.instance_id,
1899 generation: key.generation,
1900 observation: ControlObservation::ScreenState { state: merged },
1901 });
1902 }
1903 owned.next_foreground_probe =
1904 match foreground_probe_schedule(&owned.last_screen_state) {
1905 ForegroundProbeSchedule::Disarmed => None,
1906 ForegroundProbeSchedule::Armed => {
1907 Some(now + FOREGROUND_RECLASSIFY_INTERVAL)
1908 }
1909 };
1910 }
1911 Err(_) => {
1912 owned.next_foreground_probe = Some(now + FOREGROUND_RECLASSIFY_INTERVAL);
1919 }
1920 }
1921 }
1922 observations
1923 }
1924}
1925
1926const PTY_TERM_DEFAULT_KEY: &str = "TERM";
1933const PTY_TERM_DEFAULT_VALUE: &str = "xterm-256color";
1934const PTY_COLORTERM_DEFAULT_KEY: &str = "COLORTERM";
1935const PTY_COLORTERM_DEFAULT_VALUE: &str = "truecolor";
1936
1937fn with_pty_terminal_capability_defaults(mut pty_env: Vec<EnvMutation>) -> Vec<EnvMutation> {
1943 let already_mutated = |env: &[EnvMutation], key: &str| {
1944 env.iter()
1945 .any(|mutation| mutation.key.as_os_str() == OsStr::new(key))
1946 };
1947 if !already_mutated(&pty_env, PTY_TERM_DEFAULT_KEY) {
1948 pty_env.push(EnvMutation {
1949 key: OsString::from(PTY_TERM_DEFAULT_KEY),
1950 value: Some(OsString::from(PTY_TERM_DEFAULT_VALUE)),
1951 });
1952 }
1953 if !already_mutated(&pty_env, PTY_COLORTERM_DEFAULT_KEY) {
1954 pty_env.push(EnvMutation {
1955 key: OsString::from(PTY_COLORTERM_DEFAULT_KEY),
1956 value: Some(OsString::from(PTY_COLORTERM_DEFAULT_VALUE)),
1957 });
1958 }
1959 pty_env
1960}
1961
1962fn missing_session_message(key: NativeSessionKey) -> String {
1963 format!(
1964 "native session {:?}/{:?} does not exist",
1965 key.instance_id, key.generation
1966 )
1967}
1968
1969fn host_policy_for_approval_level(level: ApprovalLevel) -> HostPolicy {
1994 match level {
1995 ApprovalLevel::FullAuto => HostPolicy::Yolo,
1996 ApprovalLevel::ReadOnly => HostPolicy::ReadOnly,
1997 ApprovalLevel::Moderate | ApprovalLevel::Unmanaged => HostPolicy::Auto,
1998 }
1999}
2000
2001fn defers_permission_requests(agent_id: &AgentId, level: ApprovalLevel) -> bool {
2032 match approval_level_resolution(agent_id, level) {
2033 ApprovalLevelResolution::Supported {
2034 asks_for_permission,
2035 ..
2036 } => asks_for_permission,
2037 ApprovalLevelResolution::Unsupported => true,
2038 }
2039}
2040
2041const ACP_MODE_CONFIRMATION_TIMEOUT: Duration = Duration::from_secs(5);
2051
2052fn required_acp_mode(agent_id: &AgentId, level: ApprovalLevel) -> Result<Option<ModeId>, String> {
2069 let (acp_mode_id, carries_argv_flags) = match approval_level_resolution(agent_id, level) {
2070 ApprovalLevelResolution::Unsupported => (None, false),
2071 ApprovalLevelResolution::Supported { acp_mode_id, args, .. } => {
2072 (acp_mode_id, !args.is_empty())
2073 }
2074 };
2075 match (acp_mode_id, level) {
2076 (Some(mode_id), _) => Ok(Some(mode_id)),
2077 (None, ApprovalLevel::Unmanaged) => Ok(None),
2078 (None, _) if carries_argv_flags => Ok(None),
2087 (None, _) => Err(format!(
2088 "agent '{agent_id}' has no verified ACP mode for {level:?} and the level carries no \
2089 vendor flags to apply through argv either; refusing rather than launching at the \
2090 vendor's own (wider-authority) default"
2091 )),
2092 }
2093}
2094
2095fn acp_approval_level_args(agent_id: &AgentId, level: ApprovalLevel) -> Vec<String> {
2112 match approval_level_resolution(agent_id, level) {
2113 ApprovalLevelResolution::Supported { args, acp_mode_id: None, .. } => args,
2114 ApprovalLevelResolution::Supported { acp_mode_id: Some(_), .. } => Vec::new(),
2115 ApprovalLevelResolution::Unsupported => Vec::new(),
2116 }
2117}
2118
2119#[derive(Debug)]
2127struct ApprovalLevelNotOfferedByAgent {
2128 level: ApprovalLevel,
2129 offered: Vec<String>,
2130}
2131
2132impl std::fmt::Display for ApprovalLevelNotOfferedByAgent {
2133 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
2134 write!(
2135 formatter,
2136 "ApprovalLevelNotOfferedByAgent {{ level: {:?}, offered: {:?} }}",
2137 self.level, self.offered
2138 )
2139 }
2140}
2141
2142fn approval_level_not_offered_message(
2148 level: ApprovalLevel,
2149 mode_id: &ModeId,
2150 offered: &[SessionMode],
2151) -> Option<String> {
2152 if offered.iter().any(|mode| mode.id == mode_id.as_str()) {
2153 None
2154 } else {
2155 Some(
2156 ApprovalLevelNotOfferedByAgent {
2157 level,
2158 offered: offered.iter().map(|mode| mode.id.clone()).collect(),
2159 }
2160 .to_string(),
2161 )
2162 }
2163}
2164
2165async fn apply_acp_approval_mode(
2183 session: &AcpSession,
2184 mode_id: &ModeId,
2185 level: ApprovalLevel,
2186) -> Result<(), String> {
2187 if let Some(message) =
2188 approval_level_not_offered_message(level, mode_id, &session.available_modes())
2189 {
2190 return Err(message);
2191 }
2192
2193 let mut confirmation = session.subscribe();
2196 session
2197 .set_mode(mode_id.as_str())
2198 .await
2199 .map_err(|error| format!("session/set_mode to '{mode_id}' failed: {error}"))?;
2200
2201 let _ = tokio::time::timeout(ACP_MODE_CONFIRMATION_TIMEOUT, async {
2202 loop {
2203 match confirmation.recv().await {
2204 Ok(AgentEvent::ModeChanged { mode_id: changed }) if changed == mode_id.as_str() => {
2205 return;
2206 }
2207 Ok(_) => continue,
2208 Err(_) => return,
2209 }
2210 }
2211 })
2212 .await;
2213
2214 Ok(())
2215}
2216
2217fn validate_spawn_runtime_policy(
2218 policy: ProviderRuntimePolicy,
2219 transport: TransportKind,
2220 has_initial_prompt: bool,
2221 is_resume: bool,
2222) -> Result<(), String> {
2223 policy
2224 .validate()
2225 .map_err(|error| format!("provider runtime policy is invalid: {error}"))?;
2226 if transport != TransportKind::Acp {
2235 require_runtime_capability(policy, ProviderRuntimeCapability::RawPtyLifecycle)?;
2236 if transport != TransportKind::Pty {
2237 require_runtime_capability(policy, ProviderRuntimeCapability::SemanticReadiness)?;
2238 }
2239 if has_initial_prompt {
2240 require_runtime_capability(policy, ProviderRuntimeCapability::SemanticReadiness)?;
2241 require_runtime_capability(policy, ProviderRuntimeCapability::StructuredPrompt)?;
2242 }
2243 if is_resume && has_initial_prompt {
2244 require_runtime_capability(policy, ProviderRuntimeCapability::ProviderSessionIdentity)?;
2245 require_runtime_capability(policy, ProviderRuntimeCapability::SemanticResume)?;
2246 }
2247 }
2248 Ok(())
2249}
2250
2251
2252fn is_codex_config_c_overlay_args(arguments: &[OsString]) -> bool {
2253 if arguments.is_empty() || arguments.len() % 2 != 0 {
2254 return false;
2255 }
2256 arguments.chunks_exact(2).all(|pair| {
2257 pair[0].as_os_str() == "-c"
2258 && pair[1]
2259 .to_str()
2260 .is_some_and(|value| !value.is_empty() && value.contains('=') && !value.contains('\0'))
2261 })
2262}
2263
2264fn validate_instance_launch_arguments(
2265 agent_id: &AgentId,
2266 transport: TransportKind,
2267 arguments: &[OsString],
2268) -> Result<(), String> {
2269 if arguments.is_empty() {
2270 return Ok(());
2271 }
2272 let acp_codex_c_overlay = transport == TransportKind::Acp
2273 && agent_id.as_str() == "codex"
2274 && is_codex_config_c_overlay_args(arguments);
2275 if transport != TransportKind::Pty && !acp_codex_c_overlay {
2276 return Err(
2277 "native instance launch arguments require PTY transport (or Codex -c overlays on ACP)"
2278 .to_owned(),
2279 );
2280 }
2281 if arguments.len() > INSTANCE_LAUNCH_ARGS_MAX {
2282 return Err("native instance launch argument count exceeds its bound".to_owned());
2283 }
2284 let mut total_bytes = 0usize;
2285 for argument in arguments {
2286 let argument = argument.to_string_lossy();
2287 if argument.contains('\0') {
2288 return Err("native instance launch argument is invalid".to_owned());
2289 }
2290 if argument.len() > INSTANCE_LAUNCH_ARG_MAX_BYTES {
2291 return Err("native instance launch argument exceeds its bound".to_owned());
2292 }
2293 total_bytes = total_bytes.saturating_add(argument.len());
2294 if total_bytes > INSTANCE_LAUNCH_ARGS_TOTAL_MAX_BYTES {
2295 return Err("native instance launch argument payload exceeds its bound".to_owned());
2296 }
2297 if agent_id.as_str() == "claude"
2298 && RESERVED_CLAUDE_LAUNCH_FLAGS.iter().any(|reserved| {
2299 argument == *reserved
2300 || (reserved.starts_with("--")
2301 && argument
2302 .strip_prefix(reserved)
2303 .is_some_and(|suffix| suffix.starts_with('=')))
2304 || (reserved.len() == 2
2305 && argument
2306 .strip_prefix(reserved)
2307 .is_some_and(|suffix| {
2308 !suffix.is_empty() && !suffix.starts_with('-')
2309 }))
2310 })
2311 {
2312 return Err(
2313 "native instance launch arguments conflict with Claude session, resume, or prompt authority"
2314 .to_owned(),
2315 );
2316 }
2317 }
2318 Ok(())
2319}
2320
2321fn argument_looks_like_credential(value: &str) -> bool {
2334 const CREDENTIAL_PREFIXES: &[&str] = &["sk-", "sk_", "ghp_", "gho_", "ghs_", "xox", "bearer "];
2335 const CREDENTIAL_MARKERS: &[&str] = &[
2336 "apikey", "api_key", "api-key", "secret", "password", "passwd", "token=",
2337 ];
2338 let lower = value.to_ascii_lowercase();
2339 CREDENTIAL_PREFIXES
2340 .iter()
2341 .any(|prefix| lower.starts_with(prefix))
2342 || CREDENTIAL_MARKERS
2343 .iter()
2344 .any(|marker| lower.contains(marker))
2345 || has_prefixed_hex64_credential_shape(&lower)
2346}
2347
2348fn has_prefixed_hex64_credential_shape(lower: &str) -> bool {
2357 let Some((prefix, digest)) = lower.split_once('_') else {
2358 return false;
2359 };
2360 (2..=8).contains(&prefix.len())
2361 && prefix
2362 .bytes()
2363 .all(|byte| byte.is_ascii_lowercase() || byte.is_ascii_digit())
2364 && digest.len() == 64
2365 && digest
2366 .bytes()
2367 .all(|byte| byte.is_ascii_digit() || matches!(byte, b'a'..=b'f'))
2368}
2369
2370fn redact_provider_argument(value: &OsStr) -> String {
2374 let text = value.to_string_lossy();
2375 if argument_looks_like_credential(&text) {
2376 "[redacted-credential-shaped-argument]".to_owned()
2377 } else {
2378 text.into_owned()
2379 }
2380}
2381
2382fn redact_provider_arguments(values: &[OsString]) -> Vec<String> {
2385 values.iter().map(|value| redact_provider_argument(value)).collect()
2386}
2387
2388fn require_runtime_capability(
2389 policy: ProviderRuntimePolicy,
2390 capability: ProviderRuntimeCapability,
2391) -> Result<(), String> {
2392 if policy.admits(capability) {
2393 Ok(())
2394 } else {
2395 Err(format!(
2396 "provider runtime capability {capability:?} is not admitted"
2397 ))
2398 }
2399}
2400
2401fn should_attach_pty_provider_stream(policy: ProviderRuntimePolicy) -> bool {
2402 policy.semantic_readiness || policy.provider_session_identity
2403}
2404
2405fn should_probe_pty_identity(
2406 policy: ProviderRuntimePolicy,
2407 adapter: Option<&gate4agent_types::AdapterBinding>,
2408 authoritative_identity_present: bool,
2409 expected_adapter: &str,
2410) -> bool {
2411 !authoritative_identity_present
2412 && policy.semantic_readiness
2413 && policy.structured_prompt
2414 && policy.provider_session_identity
2415 && adapter.is_some_and(|adapter| adapter.id.as_str() == expected_adapter)
2416}
2417
2418fn prepare_fresh_pty_provider_session(
2419 adapter: Option<&gate4agent_types::AdapterBinding>,
2420 is_resume: bool,
2421 identity_permitted: bool,
2422 launch_extra_args: &mut Vec<OsString>,
2423) -> Option<ProviderSessionIdentity> {
2424 let adapter = adapter?;
2425 if is_resume || !identity_permitted || adapter.id.as_str() != "claude-code" {
2426 return None;
2427 }
2428 let identity = ProviderSessionIdentity {
2429 key: ProviderSessionKey::SessionId,
2430 id: Uuid::new_v4().to_string(),
2431 transcript_path: None,
2432 };
2433 launch_extra_args.push(OsString::from("--session-id"));
2434 launch_extra_args.push(OsString::from(&identity.id));
2435 Some(identity)
2436}
2437
2438fn drain_pty_provider(
2439 key: NativeSessionKey,
2440 provider: &mut OwnedPtyProvider,
2441 observations: &mut Vec<ObservationEnvelope>,
2442) {
2443 while let Some(event) = provider.pending_events.pop_front() {
2444 push_provider_observation(key, provider, event, observations);
2445 }
2446
2447 loop {
2448 let envelope = match provider.replay.pop_front() {
2449 Some(envelope) => Some(envelope),
2450 None => provider.receiver.try_recv().unwrap_or_default(),
2451 };
2452 let Some(envelope) = envelope else {
2453 break;
2454 };
2455 match envelope.event {
2456 PtyEvent::Output(data) => {
2457 let raw = provider.utf8.push(&data);
2458 if raw.is_empty() {
2459 continue;
2460 }
2461 if provider.semantic_events {
2462 if let Some(info) = provider.rate_limits.detect(&raw) {
2463 push_provider_observation(
2464 key,
2465 provider,
2466 rate_limit_event(info),
2467 observations,
2468 );
2469 }
2470 }
2471 let identity = provider
2472 .kimi_identity
2473 .as_mut()
2474 .and_then(|extractor| extractor.push(&raw));
2475 if let Some(identity) = identity {
2476 if !provider.provider_session_started {
2477 push_provider_observation(
2478 key,
2479 provider,
2480 ProviderEvent::SessionStarted {
2481 session_id: identity.id.clone(),
2482 model: String::new(),
2483 tools: Vec::new(),
2484 },
2485 observations,
2486 );
2487 provider.provider_session_started = true;
2488 }
2489 push_provider_observation(
2490 key,
2491 provider,
2492 ProviderEvent::SessionIdentityObserved { identity },
2493 observations,
2494 );
2495 }
2496 if provider.semantic_events {
2497 let messages = provider
2498 .pipeline
2499 .get_mut()
2500 .unwrap_or_else(|poisoned| poisoned.into_inner())
2501 .process(&raw);
2502 for message in messages {
2503 if let Some(event) = parsed_provider_event(message) {
2504 push_provider_observation(key, provider, event, observations);
2505 }
2506 }
2507 }
2508 }
2509 PtyEvent::DataGap {
2510 from_sequence,
2511 to_sequence,
2512 ..
2513 } => {
2514 provider.utf8.clear();
2515 if let Some(extractor) = &mut provider.kimi_identity {
2516 extractor.reset_stream();
2517 }
2518 provider
2519 .pipeline
2520 .get_mut()
2521 .unwrap_or_else(|poisoned| poisoned.into_inner())
2522 .clear();
2523 let missed = to_sequence
2524 .checked_sub(from_sequence)
2525 .and_then(|difference| difference.checked_add(1));
2526 if let Some((source_sequence, missed)) = missed.and_then(|missed| {
2527 reserve_provider_gap_sequence(&mut provider.next_provider_sequence, missed)
2528 .map(|source_sequence| (source_sequence, missed))
2529 }) {
2530 observations.push(ObservationEnvelope {
2531 operation_id: None,
2532 instance_id: key.instance_id,
2533 generation: key.generation,
2534 observation: ControlObservation::ProviderGap {
2535 source: provider.source.clone(),
2536 source_sequence,
2537 missed,
2538 },
2539 });
2540 }
2541 }
2542 PtyEvent::ReaderError { message } | PtyEvent::OperatorActionRequired { message } => {
2543 push_provider_observation(
2544 key,
2545 provider,
2546 ProviderEvent::Error { message },
2547 observations,
2548 );
2549 }
2550 PtyEvent::Started
2551 | PtyEvent::Resized(_)
2552 | PtyEvent::ForegroundProcess(_)
2553 | PtyEvent::SnapshotAvailable { .. }
2554 | PtyEvent::Exited { .. } => {}
2555 }
2556 }
2557}
2558
2559fn push_provider_observation(
2560 key: NativeSessionKey,
2561 provider: &mut OwnedPtyProvider,
2562 event: ProviderEvent,
2563 observations: &mut Vec<ObservationEnvelope>,
2564) {
2565 let sequence = provider.next_provider_sequence;
2566 provider.next_provider_sequence = provider.next_provider_sequence.saturating_add(1);
2567 observations.push(ObservationEnvelope {
2568 operation_id: None,
2569 instance_id: key.instance_id,
2570 generation: key.generation,
2571 observation: ControlObservation::ProviderEvent {
2572 source: provider.source.clone(),
2573 sequence,
2574 event,
2575 },
2576 });
2577}
2578
2579fn reserve_provider_gap_sequence(next_sequence: &mut u64, missed: u64) -> Option<u64> {
2580 if missed == 0 {
2581 return None;
2582 }
2583 let Some(source_sequence) = missed
2584 .checked_sub(1)
2585 .and_then(|offset| next_sequence.checked_add(offset))
2586 else {
2587 *next_sequence = u64::MAX;
2588 return None;
2589 };
2590 let Some(next) = source_sequence.checked_add(1) else {
2591 *next_sequence = u64::MAX;
2592 return None;
2593 };
2594 *next_sequence = next;
2595 Some(source_sequence)
2596}
2597
2598fn parsed_provider_event(message: ParsedMessage) -> Option<ProviderEvent> {
2599 match message.class {
2600 MessageClass::AiResponse => Some(ProviderEvent::Text {
2601 text: message.content,
2602 is_delta: message.metadata.is_partial,
2603 }),
2604 MessageClass::ThinkingIndicator => Some(ProviderEvent::Thinking {
2605 text: message.content,
2606 }),
2607 MessageClass::Error => Some(ProviderEvent::Error {
2608 message: message.content,
2609 }),
2610 MessageClass::PromptReady => Some(ProviderEvent::Ready),
2611 MessageClass::ToolApproval => Some(ProviderEvent::InteractionRequested {
2612 request_id: None,
2613 interaction_kind: ProviderInteractionKind::Approval,
2614 tool_name: message
2615 .metadata
2616 .tool_name
2617 .unwrap_or_else(|| "unknown".to_owned()),
2618 title: None,
2622 prompt: message.content,
2623 options: Vec::new(),
2624 agent_id: None,
2625 }),
2626 MessageClass::InfoMessage
2627 | MessageClass::UiElement
2628 | MessageClass::UserEcho
2629 | MessageClass::Menu
2630 | MessageClass::Raw => None,
2631 }
2632}
2633
2634fn rate_limit_event(info: gate4agent::core::types::RateLimitInfo) -> ProviderEvent {
2635 ProviderEvent::RateLimited {
2636 limit_type: provider_rate_limit_kind(info.limit_type),
2637 resets_at: info.resets_at.map(|value| value.to_rfc3339()),
2638 usage_percent: info.usage_percent.map(|value| value.to_string()),
2639 raw_message: info.raw_message,
2640 }
2641}
2642
2643fn provider_host_request_decision(decision: HostRequestDecision) -> ProviderHostRequestDecision {
2650 match decision {
2651 HostRequestDecision::Granted { by } => {
2652 ProviderHostRequestDecision::Granted { by: provider_host_decision_authority(by) }
2653 }
2654 HostRequestDecision::Denied { by } => {
2655 ProviderHostRequestDecision::Denied { by: provider_host_decision_authority(by) }
2656 }
2657 HostRequestDecision::Deferred => ProviderHostRequestDecision::Deferred,
2658 }
2659}
2660
2661fn provider_host_decision_authority(by: HostDecisionAuthority) -> ProviderHostDecisionAuthority {
2662 match by {
2663 HostDecisionAuthority::Gate => ProviderHostDecisionAuthority::Gate,
2664 HostDecisionAuthority::Policy => ProviderHostDecisionAuthority::Policy,
2665 HostDecisionAuthority::Operator => ProviderHostDecisionAuthority::Operator,
2666 HostDecisionAuthority::DeadlinePolicy => ProviderHostDecisionAuthority::DeadlinePolicy,
2667 }
2668}
2669
2670fn provider_host_request_outcome(outcome: HostRequestOutcome) -> ProviderHostRequestOutcome {
2678 match outcome {
2679 HostRequestOutcome::Executed => ProviderHostRequestOutcome::Executed,
2680 HostRequestOutcome::Failed { error } => ProviderHostRequestOutcome::Failed { error },
2681 }
2682}
2683
2684fn provider_rate_limit_kind(
2688 limit_type: gate4agent::core::types::RateLimitType,
2689) -> ProviderRateLimitKind {
2690 use gate4agent::core::types::RateLimitType;
2691 match limit_type {
2692 RateLimitType::Session => ProviderRateLimitKind::Session,
2693 RateLimitType::Daily => ProviderRateLimitKind::Daily,
2694 RateLimitType::Weekly => ProviderRateLimitKind::Weekly,
2695 RateLimitType::Unknown => ProviderRateLimitKind::Unknown,
2696 }
2697}
2698
2699fn collect_provider_map<S>(
2705 sessions: &mut BTreeMap<NativeSessionKey, OwnedProviderSession<S>>,
2706 observations: &mut Vec<ObservationEnvelope>,
2707 mode_catalogue: impl Fn(&S) -> Vec<SessionMode>,
2708) {
2709 for (key, owned) in sessions {
2710 let available_modes = mode_catalogue(&owned.session);
2711 drain_provider_stream(
2712 *key,
2713 &owned.source,
2714 &mut owned.events,
2715 &mut owned.pending_events,
2716 &mut owned.pending_provider_events,
2717 &mut owned.next_provider_sequence,
2718 Some(&mut owned.observed_exit_code),
2719 observations,
2720 &available_modes,
2721 );
2722 }
2723}
2724
2725fn push_sequenced_provider_event(
2732 key: NativeSessionKey,
2733 source: &ProviderSource,
2734 next_provider_sequence: &mut u64,
2735 observations: &mut Vec<ObservationEnvelope>,
2736 event: ProviderEvent,
2737) {
2738 let sequence = *next_provider_sequence;
2739 *next_provider_sequence = next_provider_sequence.saturating_add(1);
2740 observations.push(ObservationEnvelope {
2741 operation_id: None,
2742 instance_id: key.instance_id,
2743 generation: key.generation,
2744 observation: ControlObservation::ProviderEvent {
2745 source: source.clone(),
2746 sequence,
2747 event,
2748 },
2749 });
2750}
2751
2752fn drain_provider_stream(
2753 key: NativeSessionKey,
2754 source: &ProviderSource,
2755 events: &mut broadcast::Receiver<AgentEvent>,
2756 pending_events: &mut VecDeque<AgentEvent>,
2757 pending_provider_events: &mut VecDeque<ProviderEvent>,
2758 next_provider_sequence: &mut u64,
2759 mut observed_exit_code: Option<&mut Option<i32>>,
2760 observations: &mut Vec<ObservationEnvelope>,
2761 available_modes: &[SessionMode],
2762) {
2763 loop {
2764 if let Some(event) = pending_events.pop_front() {
2774 match event {
2775 AgentEvent::Exited { code } => {
2776 if let Some(exit_code) = observed_exit_code.as_deref_mut() {
2777 *exit_code = Some(code);
2778 }
2779 }
2780 event => {
2781 if let Some(event) = provider_event(event, available_modes) {
2782 push_sequenced_provider_event(
2783 key,
2784 source,
2785 next_provider_sequence,
2786 observations,
2787 event,
2788 );
2789 }
2790 }
2791 }
2792 continue;
2793 }
2794 if let Some(event) = pending_provider_events.pop_front() {
2795 push_sequenced_provider_event(key, source, next_provider_sequence, observations, event);
2796 continue;
2797 }
2798 match events.try_recv() {
2799 Ok(AgentEvent::Exited { code }) => {
2800 if let Some(exit_code) = observed_exit_code.as_deref_mut() {
2801 *exit_code = Some(code);
2802 }
2803 }
2804 Ok(event) => {
2805 if let Some(event) = provider_event(event, available_modes) {
2806 push_sequenced_provider_event(
2807 key,
2808 source,
2809 next_provider_sequence,
2810 observations,
2811 event,
2812 );
2813 }
2814 }
2815 Err(broadcast::error::TryRecvError::Lagged(missed)) => {
2816 if let Some(source_sequence) =
2817 reserve_provider_gap_sequence(next_provider_sequence, missed)
2818 {
2819 observations.push(ObservationEnvelope {
2820 operation_id: None,
2821 instance_id: key.instance_id,
2822 generation: key.generation,
2823 observation: ControlObservation::ProviderGap {
2824 source: source.clone(),
2825 source_sequence,
2826 missed,
2827 },
2828 });
2829 }
2830 }
2831 Err(broadcast::error::TryRecvError::Empty | broadcast::error::TryRecvError::Closed) => {
2832 break
2833 }
2834 }
2835 }
2836}
2837
2838fn collect_provider_exits<S>(
2839 sessions: &mut BTreeMap<NativeSessionKey, OwnedProviderSession<S>>,
2840 observations: &mut Vec<ObservationEnvelope>,
2841 finished: impl Fn(&S) -> bool,
2842) {
2843 let completed: Vec<_> = sessions
2844 .iter()
2845 .filter_map(|(key, owned)| {
2846 (owned.observed_exit_code.is_some() || finished(&owned.session)).then_some(*key)
2847 })
2848 .collect();
2849 for key in completed {
2850 let owned = sessions
2851 .remove(&key)
2852 .expect("completed provider key came from the owned session map");
2853 observations.push(ObservationEnvelope {
2854 operation_id: None,
2855 instance_id: key.instance_id,
2856 generation: key.generation,
2857 observation: ControlObservation::ProcessExited {
2858 exit_code: owned.observed_exit_code,
2859 final_terminal: None,
2860 },
2861 });
2862 }
2863}
2864
2865const RPC_REQUEST_ID_NUMBER_PREFIX: &str = "number:";
2871const RPC_REQUEST_ID_STRING_PREFIX: &str = "string:";
2872
2873fn encode_rpc_request_id(id: &RpcId) -> String {
2885 match id {
2886 RpcId::Number(number) => format!("{RPC_REQUEST_ID_NUMBER_PREFIX}{number}"),
2887 RpcId::String(value) => format!("{RPC_REQUEST_ID_STRING_PREFIX}{value}"),
2888 }
2889}
2890
2891fn decode_rpc_request_id(raw: &str) -> Option<RpcId> {
2896 if let Some(number) = raw.strip_prefix(RPC_REQUEST_ID_NUMBER_PREFIX) {
2897 return number.parse::<u64>().ok().map(RpcId::Number);
2898 }
2899 raw.strip_prefix(RPC_REQUEST_ID_STRING_PREFIX)
2900 .map(|value| RpcId::String(value.to_owned()))
2901}
2902
2903fn resolve_acp_interaction_observation(
2916 session: Option<&AcpSession>,
2917 target: ProviderInteractionTarget,
2918 response: ProviderInteractionResponse,
2919) -> ControlObservation {
2920 let interaction_id = target.interaction_id;
2921 let fail = |message: String| ControlObservation::InteractionResolutionFailed {
2922 interaction_id,
2923 message,
2924 };
2925 let Some(session) = session else {
2926 return fail("semantic interaction resolution requires an ACP session".to_owned());
2927 };
2928 let Some(provider_request_id) = target.provider_request_id.as_deref() else {
2929 return fail("ACP interaction resolution requires a provider request id".to_owned());
2930 };
2931 let Some(id) = decode_rpc_request_id(provider_request_id) else {
2932 return fail(format!(
2933 "ACP interaction resolution request id {provider_request_id:?} is not a valid JSON-RPC id"
2934 ));
2935 };
2936 match resolve_acp_permission_interaction(session, &id, response) {
2937 Ok(()) => ControlObservation::InteractionResolutionCompleted { interaction_id },
2938 Err(message) => fail(message),
2939 }
2940}
2941
2942fn resolve_acp_permission_interaction(
2966 session: &AcpSession,
2967 id: &RpcId,
2968 response: ProviderInteractionResponse,
2969) -> Result<(), String> {
2970 let choice = match response {
2971 ProviderInteractionResponse::ApproveOnce => OperatorPermissionChoice::Approve,
2972 ProviderInteractionResponse::Deny => OperatorPermissionChoice::Reject,
2973 ProviderInteractionResponse::Answer { .. } => {
2974 return Err("ACP permission interactions do not accept free-text answers".to_owned())
2975 }
2976 };
2977 session.resolve_pending_request_as(id, choice).map_err(|error| error.to_string())
2978}
2979
2980async fn set_acp_session_mode_observation(
2993 session: Option<&AcpSession>,
2994 mode_id: String,
2995) -> ControlObservation {
2996 let Some(session) = session else {
2997 return ControlObservation::SessionModeSetFailed {
2998 message: "session mode switch requires an ACP session".to_owned(),
2999 };
3000 };
3001 let catalogue = session.available_modes();
3002 if catalogue.is_empty() {
3003 return ControlObservation::SessionModeSetFailed {
3004 message: "the provider announces no session-mode capability".to_owned(),
3005 };
3006 }
3007 if !catalogue.iter().any(|mode| mode.id == mode_id) {
3008 return ControlObservation::SessionModeSetFailed {
3009 message: format!("the provider never offered mode {mode_id:?}"),
3010 };
3011 }
3012 match session.set_mode(&mode_id).await {
3013 Ok(()) => ControlObservation::SessionModeSet { mode_id },
3014 Err(error) => ControlObservation::SessionModeSetFailed {
3015 message: error.to_string(),
3016 },
3017 }
3018}
3019
3020async fn set_acp_session_config_option_observation(
3036 session: Option<&AcpSession>,
3037 option_id: String,
3038 value_json: String,
3039) -> ControlObservation {
3040 let Some(session) = session else {
3041 return ControlObservation::SessionConfigOptionSetFailed {
3042 message: "session config-option switch requires an ACP session".to_owned(),
3043 };
3044 };
3045 let catalogue = session.config_options();
3046 if catalogue.is_empty() {
3047 return ControlObservation::SessionConfigOptionSetFailed {
3048 message: "the provider announces no session config-option capability".to_owned(),
3049 };
3050 }
3051 if !catalogue.iter().any(|option| option.id == option_id) {
3052 return ControlObservation::SessionConfigOptionSetFailed {
3053 message: format!("the provider never offered config option {option_id:?}"),
3054 };
3055 }
3056 let value = match value_json.parse() {
3057 Ok(value) => value,
3058 Err(error) => {
3059 return ControlObservation::SessionConfigOptionSetFailed {
3060 message: format!("session config option value is not valid JSON: {error}"),
3061 }
3062 }
3063 };
3064 match session.set_config_option(&option_id, value).await {
3065 Ok(()) => ControlObservation::SessionConfigOptionSet { option_id },
3066 Err(error) => ControlObservation::SessionConfigOptionSetFailed {
3067 message: error.to_string(),
3068 },
3069 }
3070}
3071
3072fn set_acp_session_model_observation(
3081 session: Option<&AcpSession>,
3082 model_id: String,
3083) -> ControlObservation {
3084 if session.is_none() {
3085 return ControlObservation::SessionModelSetFailed {
3086 message: "session model switch requires an ACP session".to_owned(),
3087 };
3088 }
3089 ControlObservation::SessionModelSetFailed {
3090 message: format!(
3091 "this provider offers no host-invocable model-switch method (requested model {model_id:?})"
3092 ),
3093 }
3094}
3095
3096fn map_provider_stop_reason(reason: StopReason) -> ProviderStopReason {
3109 match reason {
3110 StopReason::EndTurn => ProviderStopReason::EndTurn,
3111 StopReason::MaxTokens => ProviderStopReason::MaxTokens,
3112 StopReason::MaxTurnRequests => ProviderStopReason::MaxTurnRequests,
3113 StopReason::Refusal => ProviderStopReason::Refusal,
3114 StopReason::Cancelled => ProviderStopReason::Cancelled,
3115 StopReason::Other(value) => ProviderStopReason::Other { value },
3116 StopReason::ProviderError { code, message, vendor_code } => {
3117 ProviderStopReason::ProviderError { code, message, vendor_code }
3118 }
3119 }
3120}
3121
3122fn provider_event(event: AgentEvent, available_modes: &[SessionMode]) -> Option<ProviderEvent> {
3123 match event {
3124 AgentEvent::SessionStart {
3125 session_id,
3126 model,
3127 tools,
3128 } => Some(ProviderEvent::SessionStarted {
3129 session_id,
3130 model,
3131 tools,
3132 }),
3133 AgentEvent::Text { text, is_delta } => Some(ProviderEvent::Text { text, is_delta }),
3134 AgentEvent::Thinking { text } => Some(ProviderEvent::Thinking { text }),
3135 AgentEvent::ToolStart { id, name, input } => Some(ProviderEvent::ToolStarted {
3136 id,
3137 name,
3138 input_json: input.to_string(),
3139 agent_id: None,
3140 }),
3141 AgentEvent::ToolResult {
3142 id,
3143 output,
3144 is_error,
3145 duration_ms,
3146 non_execution_kind,
3147 } => Some(ProviderEvent::ToolCompleted {
3148 id,
3149 output,
3150 is_error,
3151 duration_ms,
3152 agent_id: None,
3153 non_execution_kind,
3154 }),
3155 AgentEvent::TurnComplete {
3156 input_tokens,
3157 output_tokens,
3158 cache_read_tokens,
3159 cache_write_tokens,
3160 reasoning_tokens,
3161 context_window,
3162 is_cumulative,
3163 } => Some(ProviderEvent::TurnCompleted {
3164 usage: TokenUsage {
3165 input_tokens,
3166 output_tokens,
3167 cache_read_tokens,
3168 cache_write_tokens,
3169 reasoning_tokens,
3170 context_window,
3171 },
3172 is_cumulative,
3173 }),
3174 AgentEvent::ContextWindowUsage { usage } => {
3175 Some(ProviderEvent::ContextWindowUsage {
3176 usage: ProviderContextWindowUsage {
3177 uncached_input_tokens: usage.uncached_input_tokens,
3178 cache_read_tokens: usage.cache_read_tokens,
3179 cache_write_tokens: usage.cache_write_tokens,
3180 output_tokens: usage.output_tokens,
3181 unattributed_tokens: usage.unattributed_tokens,
3182 used_tokens: usage.used_tokens,
3183 capacity_tokens: usage.capacity_tokens,
3184 },
3185 })
3186 }
3187 AgentEvent::SessionEnd {
3188 result,
3189 cost_usd,
3190 is_error,
3191 stop_reason,
3192 } => Some(ProviderEvent::SessionEnded {
3193 result,
3194 cost_usd: cost_usd.map(|cost| cost.to_string()),
3195 is_error,
3196 stop_reason: stop_reason.map(map_provider_stop_reason),
3197 }),
3198 AgentEvent::Error { message } => Some(ProviderEvent::Error { message }),
3199 AgentEvent::TurnInterrupted { reason } => {
3209 tracing::warn!(reason = %reason, "acp turn interrupted");
3210 Some(ProviderEvent::TurnInterrupted)
3211 }
3212 AgentEvent::PtyParsed(message) => parsed_provider_event(message),
3213 AgentEvent::PtyReady => Some(ProviderEvent::Ready),
3214 AgentEvent::PtyToolApproval {
3215 tool_name,
3216 description,
3217 } => Some(ProviderEvent::InteractionRequested {
3218 request_id: None,
3219 interaction_kind: ProviderInteractionKind::Approval,
3220 tool_name,
3221 title: None,
3225 prompt: description.unwrap_or_default(),
3226 options: Vec::new(),
3227 agent_id: None,
3228 }),
3229 AgentEvent::RateLimit(info) => Some(rate_limit_event(info)),
3230 AgentEvent::RpcIncomingRequest {
3244 id,
3245 method,
3246 params,
3247 decision,
3248 reason: _,
3252 outcome: _,
3256 } if method == "session/request_permission"
3257 && matches!(decision, HostRequestDecision::Deferred) =>
3258 {
3259 let tool_call = params.as_ref().and_then(|params| params.get("toolCall"));
3260 let field = |name: &str| {
3261 tool_call
3262 .and_then(|tool_call| tool_call.get(name))
3263 .and_then(|value| value.as_str())
3264 .filter(|value| !value.is_empty())
3265 .map(str::to_owned)
3266 };
3267 let title = field("title");
3276 let options: Vec<ProviderInteractionOption> = params
3288 .as_ref()
3289 .and_then(|params| params.get("options"))
3290 .and_then(|options| options.as_array())
3291 .map(|options| {
3292 options
3293 .iter()
3294 .filter_map(|option| {
3295 let option_id = option.get("optionId")?.as_str()?.to_owned();
3296 let name = option.get("name")?.as_str()?.to_owned();
3297 let kind = option.get("kind")?.as_str()?.to_owned();
3298 Some(ProviderInteractionOption { option_id, name, kind })
3299 })
3300 .collect()
3301 })
3302 .unwrap_or_default();
3303 Some(ProviderEvent::InteractionRequested {
3304 request_id: Some(encode_rpc_request_id(&id)),
3305 interaction_kind: ProviderInteractionKind::Approval,
3306 tool_name: field("kind").unwrap_or(method),
3307 title,
3308 prompt: field("title").unwrap_or_default(),
3309 options,
3310 agent_id: None,
3311 })
3312 }
3313 AgentEvent::RpcIncomingRequest {
3314 id: _,
3315 method,
3316 params,
3317 decision,
3318 outcome,
3327 reason,
3328 } => Some(ProviderEvent::HostRequestObserved {
3329 method,
3330 params_json: params.map(|value| value.to_string()).unwrap_or_default(),
3331 decision: provider_host_request_decision(decision),
3332 outcome: provider_host_request_outcome(outcome),
3333 reason,
3334 }),
3335 AgentEvent::RpcNotification { method, params } => {
3336 Some(ProviderEvent::UnrecognizedNotification {
3337 method,
3338 payload_json: params.to_string(),
3339 })
3340 }
3341 AgentEvent::Started { .. } | AgentEvent::Exited { .. } | AgentEvent::PtyRaw { .. } => None,
3342 AgentEvent::UserMessage { text, is_delta } => {
3343 Some(ProviderEvent::UserMessage { text, is_delta })
3344 }
3345 AgentEvent::Plan { steps } => Some(ProviderEvent::Plan {
3346 steps: steps.into_iter().map(provider_plan_step).collect(),
3347 }),
3348 AgentEvent::AvailableCommandsUpdate { commands } => {
3349 Some(ProviderEvent::AvailableCommandsUpdated {
3350 commands: commands
3351 .into_iter()
3352 .map(|command| ProviderAvailableCommand {
3353 name: command.name,
3354 description: command.description,
3355 input_hint: command.input_hint,
3356 })
3357 .collect(),
3358 })
3359 }
3360 AgentEvent::ModeChanged { mode_id } => Some(ProviderEvent::ModeChanged {
3361 mode_id,
3362 available: available_modes.iter().cloned().map(provider_mode_info).collect(),
3363 }),
3364 AgentEvent::SessionInfoUpdate { title } => {
3365 Some(ProviderEvent::SessionInfoUpdated { title })
3366 }
3367 AgentEvent::UsageUpdate {
3368 used_tokens,
3369 context_window,
3370 cost_amount,
3371 cost_currency,
3372 } => Some(ProviderEvent::UsageUpdated {
3373 used_tokens,
3374 context_window,
3375 cost_amount: cost_amount.map(|amount| amount.to_string()),
3376 cost_currency,
3377 }),
3378 AgentEvent::ConfigOptionsUpdate { options } => Some(ProviderEvent::ConfigOptionsUpdated {
3379 options: options.into_iter().map(provider_config_option).collect(),
3380 }),
3381 AgentEvent::ModelsUpdate { .. }
3388 | AgentEvent::ProviderModelChanged { .. }
3389 | AgentEvent::SettingsUpdate { .. }
3390 | AgentEvent::HookExecutionUpdate { .. }
3391 | AgentEvent::McpServersUpdate { .. }
3392 | AgentEvent::McpInitProgress { .. }
3393 | AgentEvent::McpInitialized { .. }
3394 | AgentEvent::AnnouncementsUpdate { .. } => None,
3395 }
3396}
3397
3398fn provider_plan_priority(priority: gate4agent::core::types::PlanStepPriority) -> ProviderPlanPriority {
3402 use gate4agent::core::types::PlanStepPriority;
3403 match priority {
3404 PlanStepPriority::High => ProviderPlanPriority::High,
3405 PlanStepPriority::Medium => ProviderPlanPriority::Medium,
3406 PlanStepPriority::Low => ProviderPlanPriority::Low,
3407 }
3408}
3409
3410fn provider_plan_status(status: gate4agent::core::types::PlanStepStatus) -> ProviderPlanStatus {
3412 use gate4agent::core::types::PlanStepStatus;
3413 match status {
3414 PlanStepStatus::Pending => ProviderPlanStatus::Pending,
3415 PlanStepStatus::InProgress => ProviderPlanStatus::InProgress,
3416 PlanStepStatus::Completed => ProviderPlanStatus::Completed,
3417 }
3418}
3419
3420fn provider_plan_step(step: gate4agent::core::types::PlanStep) -> ProviderPlanStep {
3421 ProviderPlanStep {
3422 content: step.content,
3423 priority: provider_plan_priority(step.priority),
3424 status: provider_plan_status(step.status),
3425 }
3426}
3427
3428fn provider_mode_info(mode: SessionMode) -> ProviderModeInfo {
3432 ProviderModeInfo {
3433 id: mode.id,
3434 name: mode.name,
3435 description: mode.description,
3436 }
3437}
3438
3439fn provider_config_option_kind(
3442 kind: gate4agent::core::types::ConfigOptionKind,
3443) -> ProviderConfigOptionKind {
3444 use gate4agent::core::types::ConfigOptionKind;
3445 match kind {
3446 ConfigOptionKind::Select => ProviderConfigOptionKind::Select,
3447 ConfigOptionKind::Boolean => ProviderConfigOptionKind::Boolean,
3448 ConfigOptionKind::Unknown => ProviderConfigOptionKind::Unknown,
3449 }
3450}
3451
3452fn provider_config_option(option: gate4agent::core::types::ConfigOptionInfo) -> ProviderConfigOption {
3458 ProviderConfigOption {
3459 id: option.id,
3460 name: option.name,
3461 description: option.description,
3462 category: option.category,
3463 kind: provider_config_option_kind(option.kind),
3464 value_json: option.value.to_string(),
3465 choices: option
3466 .choices
3467 .into_iter()
3468 .map(|choice| ProviderConfigChoice {
3469 value_json: choice.value.to_string(),
3470 label: choice.label,
3471 })
3472 .collect(),
3473 }
3474}
3475
3476fn builtin_legacy_adapter_runtimes() -> AdapterRuntimeRegistry<CliTool> {
3477 let definitions = [
3483 ("claude-code", CliTool::ClaudeCode),
3484 ("codex", CliTool::Codex),
3485 ("kimi", CliTool::KimiCode),
3486 ("grok", CliTool::Grok),
3487 ];
3488 let mut runtimes = AdapterRuntimeRegistry::default();
3489 for (id, tool) in definitions {
3490 for family in [
3491 AdapterFamily::PtySemantic,
3492 AdapterFamily::Pipe,
3493 AdapterFamily::Acp,
3494 ] {
3495 let Some(binding) = builtin_adapter_registry().binding(family, id) else {
3496 continue;
3497 };
3498 runtimes
3499 .insert(family, binding.clone(), tool)
3500 .expect("built-in native adapter runtime must be unique");
3501 }
3502 }
3503
3504 runtimes
3505}
3506
3507fn terminal_frame(snapshot: PtyTerminalSnapshot, screen_state: PtyScreenState) -> TerminalFrame {
3508 TerminalFrame {
3509 sequence: snapshot.sequence,
3510 size: TerminalSize {
3511 rows: snapshot.size.rows,
3512 columns: snapshot.size.cols,
3513 },
3514 cursor_row: snapshot.cursor.0,
3515 cursor_column: snapshot.cursor.1,
3516 contents: snapshot.contents,
3517 formatted: snapshot.formatted,
3518 scrollback_formatted: snapshot.scrollback_formatted,
3519 alternate_screen: snapshot.alternate_screen,
3520 mouse_protocol_enabled: snapshot.mouse_protocol_enabled,
3521 mouse_protocol_encoding: match snapshot.mouse_protocol_encoding {
3522 PtyMouseProtocolEncoding::Default => TerminalMouseProtocolEncoding::Default,
3523 PtyMouseProtocolEncoding::Utf8 => TerminalMouseProtocolEncoding::Utf8,
3524 PtyMouseProtocolEncoding::Sgr => TerminalMouseProtocolEncoding::Sgr,
3525 },
3526 produced_at_unix_ms: snapshot.produced_at_unix_ms,
3529 screen_state,
3530 bracketed_paste: Some(snapshot.bracketed_paste),
3535 }
3536}
3537
3538fn terminal_state_capture_should_skip<E>(sequence: Result<u64, E>, last_captured: u64) -> bool {
3546 matches!(sequence, Ok(sequence) if sequence <= last_captured)
3547}
3548
3549fn terminal_frame_byte_len(frame: &TerminalFrame) -> u64 {
3555 let scrollback_bytes: usize = frame.scrollback_formatted.iter().map(Vec::len).sum();
3556 (frame.formatted.len() + scrollback_bytes) as u64
3557}
3558
3559fn canonical_foreground(
3560 agent_id: AgentId,
3561 observation: &PtyForegroundObservation,
3562) -> ForegroundProcess {
3563 let kind = if observation.readiness.process_name.as_deref() == Some(agent_id.as_str()) {
3564 ForegroundProcessKind::Agent { agent_id }
3565 } else if observation.readiness.is_shell {
3566 ForegroundProcessKind::Shell
3567 } else {
3568 ForegroundProcessKind::Other
3569 };
3570 ForegroundProcess {
3571 root_process_id: observation.root_pid,
3572 process_id: observation.observed_pid,
3573 process_name: observation.observed_process.clone(),
3574 kind,
3575 }
3576}
3577
3578fn completion_observation(
3579 operation_id: OperationId,
3580 instance_id: AgentInstanceId,
3581 generation: SessionGeneration,
3582 observation: ControlObservation,
3583) -> ObservationEnvelope {
3584 ObservationEnvelope {
3585 operation_id: Some(operation_id),
3586 instance_id,
3587 generation,
3588 observation,
3589 }
3590}
3591
3592async fn wait_for_readiness(
3593 session: &PtySession,
3594 spec: &AgentSpec,
3595 intent: ReadinessIntent,
3596 detect_startup_gates: bool,
3597) -> Result<ReadinessPermit, String> {
3598 let started = Instant::now();
3599 let (terminal, attachment) = attach_readiness_boundary(session)?;
3600 let mut receiver = attachment.receiver;
3601 let mut tracker = ReadinessTracker::new(spec, RuntimePlatform::current(), intent);
3602 let mut diagnostics = ReadinessDiagnostics {
3603 draft_signal: Some(spec.readiness.draft_signal),
3604 ..ReadinessDiagnostics::default()
3605 };
3606 seed_readiness_from_terminal(
3607 &terminal,
3608 &mut tracker,
3609 &mut diagnostics,
3610 elapsed_ms(started),
3611 )?;
3612 let interval = Duration::from_millis(spec.readiness.poll_interval_ms.max(1));
3613 let mut next_probe = Instant::now();
3614
3615 for event in attachment.replay {
3616 observe_readiness_event(
3617 &mut tracker,
3618 &mut diagnostics,
3619 event.event,
3620 elapsed_ms(started),
3621 )?;
3622 if detect_startup_gates {
3623 ensure_no_readiness_operator_gate(&diagnostics, spec)?;
3624 ensure_no_startup_operator_gate(session, spec)?;
3625 }
3626 }
3627
3628 loop {
3629 if detect_startup_gates {
3630 ensure_no_startup_operator_gate(session, spec)?;
3631 }
3632 if Instant::now() >= next_probe {
3633 let foreground = session
3634 .observe_foreground()
3635 .await
3636 .map_err(|error| error.to_string())?;
3637 diagnostics.observe_foreground(&foreground.readiness);
3638 tracker.observe_foreground(&foreground.readiness, elapsed_ms(started));
3639 next_probe = Instant::now() + interval;
3640 }
3641 tracker.poll(elapsed_ms(started));
3642 if readiness_complete(tracker.status(), &diagnostics)? {
3643 if detect_startup_gates {
3644 ensure_no_startup_operator_gate(session, spec)?;
3645 }
3646 return tracker
3647 .into_permit()
3648 .ok_or_else(|| "ready tracker did not issue a permit".to_owned());
3649 }
3650
3651 let wait = next_probe.saturating_duration_since(Instant::now());
3652 match tokio::time::timeout(wait, receiver.recv()).await {
3653 Ok(Ok(event)) => {
3654 if matches!(&event.event, PtyEvent::DataGap { .. }) {
3655 let (terminal, attachment) = attach_readiness_boundary(session)?;
3656 receiver = attachment.receiver;
3657 tracker = ReadinessTracker::new(
3658 spec,
3659 RuntimePlatform::current(),
3660 intent,
3661 );
3662 diagnostics = ReadinessDiagnostics {
3663 draft_signal: Some(spec.readiness.draft_signal),
3664 ..ReadinessDiagnostics::default()
3665 };
3666 seed_readiness_from_terminal(
3667 &terminal,
3668 &mut tracker,
3669 &mut diagnostics,
3670 elapsed_ms(started),
3671 )?;
3672 for retained in attachment.replay {
3673 observe_readiness_event(
3674 &mut tracker,
3675 &mut diagnostics,
3676 retained.event,
3677 elapsed_ms(started),
3678 )?;
3679 }
3680 } else {
3681 observe_readiness_event(
3682 &mut tracker,
3683 &mut diagnostics,
3684 event.event,
3685 elapsed_ms(started),
3686 )?;
3687 }
3688 if detect_startup_gates {
3689 ensure_no_readiness_operator_gate(&diagnostics, spec)?;
3690 ensure_no_startup_operator_gate(session, spec)?;
3691 }
3692 }
3693 Ok(Err(error)) => return Err(error.to_string()),
3694 Err(_) => {
3695 tracker.poll(elapsed_ms(started));
3696 }
3697 }
3698 if readiness_complete(tracker.status(), &diagnostics)? {
3699 if detect_startup_gates {
3700 ensure_no_startup_operator_gate(session, spec)?;
3701 }
3702 return tracker
3703 .into_permit()
3704 .ok_or_else(|| "ready tracker did not issue a permit".to_owned());
3705 }
3706 }
3707}
3708
3709fn attach_readiness_boundary(
3710 session: &PtySession,
3711) -> Result<(PtyTerminalSnapshot, PtyAttachment), String> {
3712 let terminal = session.terminal_state().map_err(|error| error.to_string())?;
3713 let cursor = PtyReplayCursor {
3714 provider_revision: terminal.provider_revision.clone(),
3715 generation: terminal.generation,
3716 next_sequence: terminal.sequence.saturating_add(1).max(1),
3717 };
3718 let attachment = session
3719 .attach_events(cursor)
3720 .map_err(|error| error.to_string())?;
3721 Ok((terminal, attachment))
3722}
3723
3724fn seed_readiness_from_terminal(
3725 terminal: &PtyTerminalSnapshot,
3726 tracker: &mut ReadinessTracker<'_>,
3727 diagnostics: &mut ReadinessDiagnostics,
3728 elapsed_ms: u64,
3729) -> Result<(), String> {
3730 const ENABLE_BRACKETED_PASTE: &[u8] = b"\x1b[?2004h";
3731 if terminal.bracketed_paste {
3732 diagnostics.observe_output(ENABLE_BRACKETED_PASTE);
3733 tracker.observe_output(ENABLE_BRACKETED_PASTE, elapsed_ms);
3734 }
3735 if !terminal.formatted.is_empty() {
3736 diagnostics.observe_output(&terminal.formatted);
3737 tracker.observe_output(&terminal.formatted, elapsed_ms);
3738 }
3739 Ok(())
3740}
3741
3742const STARTUP_GATE_SETTLE_MS: u64 = 350;
3743const STARTUP_GATE_POLL_MS: u64 = 25;
3744const FOREGROUND_RECLASSIFY_INTERVAL: Duration = Duration::from_secs(1);
3758const CODEX_PASTE_ENTER_SUPPRESSION_MS: u64 = 120;
3762const CODEX_POST_RENDER_MARGIN_MS: u64 = 30;
3763const CODEX_SESSION_STATUS_PROBE_TIMEOUT_MS: u64 = 5_000;
3764const KIMI_SESSION_STATUS_PROBE_TIMEOUT_MS: u64 = 5_000;
3765
3766async fn probe_fresh_codex_session_identity(
3767 session: &PtySession,
3768 spec: &AgentSpec,
3769) -> Result<Option<ProviderSessionIdentity>, String> {
3770 let permit = match wait_for_readiness(session, spec, ReadinessIntent::DraftPaste, true).await {
3771 Ok(permit) => permit,
3772 Err(_) => return Ok(None),
3773 };
3774 if wait_for_startup_operator_gate(session, spec).await.is_err() {
3775 return Ok(None);
3776 }
3777 let baseline = session
3778 .terminal_state()
3779 .map_err(|error| error.to_string())?;
3780 let mut extractor = CodexPtySessionIdentityExtractor::default();
3781 session
3782 .send_input_action(
3783 InputAction::AgentCommand(AgentCommand {
3784 agent_id: spec.id.clone(),
3785 name: "status".to_owned(),
3786 arguments: Vec::new(),
3787 }),
3788 permit,
3789 )
3790 .await
3791 .map_err(|error| format!("Codex /status identity probe failed: {error}"))?;
3792
3793 let deadline = Instant::now()
3794 + Duration::from_millis(
3795 spec.readiness
3796 .timeout_ms
3797 .min(CODEX_SESSION_STATUS_PROBE_TIMEOUT_MS)
3798 .max(1),
3799 );
3800 loop {
3801 let snapshot = session
3802 .terminal_state()
3803 .map_err(|error| error.to_string())?;
3804 if snapshot.sequence > baseline.sequence {
3805 if let Some(identity) = extractor.observe_screen(&snapshot.contents) {
3806 return Ok(Some(identity));
3807 }
3808 }
3809 if Instant::now() >= deadline {
3810 return Ok(None);
3811 }
3812 tokio::time::sleep(Duration::from_millis(STARTUP_GATE_POLL_MS)).await;
3813 }
3814}
3815
3816async fn probe_fresh_kimi_session_identity(
3817 session: &PtySession,
3818 spec: &AgentSpec,
3819) -> Result<Option<ProviderSessionIdentity>, String> {
3820 let permit = match wait_for_readiness(
3821 session,
3822 spec,
3823 ReadinessIntent::FollowupPrompt,
3824 true,
3825 )
3826 .await
3827 {
3828 Ok(permit) => permit,
3829 Err(_) => return Ok(None),
3830 };
3831 if wait_for_startup_operator_gate(session, spec).await.is_err() {
3832 return Ok(None);
3833 }
3834 let baseline = session
3835 .terminal_state()
3836 .map_err(|error| error.to_string())?;
3837 let mut extractor = KimiPtySessionIdentityExtractor::default();
3838 if let Some(identity) = extractor.observe_screen(&baseline.contents) {
3839 return Ok(Some(identity));
3840 }
3841 session
3842 .send_input_action(
3843 InputAction::SubmitPrompt(PromptPayload {
3844 text: "/status".to_owned(),
3845 framing: PromptFraming::BracketedPaste,
3846 }),
3847 permit,
3848 )
3849 .await
3850 .map_err(|error| format!("Kimi /status identity probe failed: {error}"))?;
3851
3852 let deadline = Instant::now()
3853 + Duration::from_millis(
3854 spec.readiness
3855 .timeout_ms
3856 .min(KIMI_SESSION_STATUS_PROBE_TIMEOUT_MS)
3857 .max(1),
3858 );
3859 loop {
3860 let snapshot = session
3861 .terminal_state()
3862 .map_err(|error| error.to_string())?;
3863 if snapshot.sequence > baseline.sequence {
3864 if let Some(identity) = extractor.observe_screen(&snapshot.contents) {
3865 return Ok(Some(identity));
3866 }
3867 }
3868 if Instant::now() >= deadline {
3869 return Ok(None);
3870 }
3871 tokio::time::sleep(Duration::from_millis(STARTUP_GATE_POLL_MS)).await;
3872 }
3873}
3874
3875async fn deliver_pending_initial_prompt(
3876 session: &mut PtySession,
3877 spec: &AgentSpec,
3878) -> Result<(), String> {
3879 if session.pending_followup_prompt().is_none() {
3880 return Ok(());
3881 }
3882 let mut permit =
3883 wait_for_readiness(session, spec, ReadinessIntent::FollowupPrompt, true).await?;
3884 wait_for_startup_operator_gate(session, spec).await?;
3885 let render_confirmed_submit = spec.id.as_str() == "claude"
3886 || (RuntimePlatform::current() == RuntimePlatform::Windows
3887 && spec.id.as_str() == "codex");
3888 if render_confirmed_submit {
3889 let prompt = session
3890 .pending_followup_prompt()
3891 .ok_or_else(|| "deferred initial prompt disappeared before paste".to_owned())?
3892 .to_owned();
3893 let baseline = session
3894 .terminal_state()
3895 .map_err(|error| error.to_string())?;
3896 let inserted = session
3897 .insert_pending_followup_prompt(PromptFraming::BracketedPaste, &permit)
3898 .await
3899 .map_err(|error| error.to_string())?;
3900 if !inserted {
3901 return Err("deferred initial prompt was not inserted".to_owned());
3902 }
3903 wait_for_prompt_render(
3904 session,
3905 spec,
3906 &prompt,
3907 &baseline,
3908 spec.readiness.timeout_ms,
3909 )
3910 .await?;
3911 if RuntimePlatform::current() == RuntimePlatform::Windows && spec.id.as_str() == "codex" {
3912 tokio::time::sleep(Duration::from_millis(
3913 CODEX_PASTE_ENTER_SUPPRESSION_MS + CODEX_POST_RENDER_MARGIN_MS,
3914 ))
3915 .await;
3916 }
3917 permit = wait_for_readiness(session, spec, ReadinessIntent::FollowupPrompt, true).await?;
3918 wait_for_startup_operator_gate(session, spec).await?;
3919 }
3920 ensure_no_startup_operator_gate(session, spec)?;
3921 let submitted = session
3922 .submit_pending_followup(PromptFraming::BracketedPaste, permit)
3923 .await
3924 .map_err(|error| error.to_string())?;
3925 if !submitted {
3926 return Err("deferred initial prompt disappeared before delivery".to_owned());
3927 }
3928 Ok(())
3929}
3930
3931async fn wait_for_prompt_render(
3932 session: &PtySession,
3933 spec: &AgentSpec,
3934 prompt: &str,
3935 baseline: &PtyTerminalSnapshot,
3936 timeout_ms: u64,
3937) -> Result<(), String> {
3938 let probe = prompt_render_probe(&gate4agent_types::sanitize_prompt_text(prompt));
3939 let deadline = Instant::now() + Duration::from_millis(timeout_ms.max(1));
3940 loop {
3941 let snapshot = session.terminal_state().map_err(|error| error.to_string())?;
3942 if let Some(gate) = startup_operator_gate(&snapshot.contents) {
3943 return Err(startup_operator_error(spec, &gate));
3944 }
3945 if prompt_rendered(&snapshot, baseline, &probe) {
3946 return Ok(());
3947 }
3948 if Instant::now() >= deadline {
3949 let tail = terminal_tail(&snapshot.contents);
3950 let compact_tail = compact_alphanumeric(&tail);
3951 let baseline_tail = terminal_tail(&baseline.contents).to_ascii_lowercase();
3952 let normalized_tail = tail.to_ascii_lowercase();
3953 let retained = render_ack_event_summary(session, baseline.sequence);
3954 let foreground = session.observe_foreground().await.ok();
3955 let foreground_name = foreground
3956 .as_ref()
3957 .and_then(|observation| observation.readiness.process_name.as_deref())
3958 .map(safe_process_label)
3959 .unwrap_or_else(|| "none".to_owned());
3960 let child_live = retained.exit_code.is_none() && foreground.is_some();
3961 return Err(format!(
3962 "agent '{}' initial prompt paste was not rendered before submit; Enter was not sent (baseline_sequence={} current_sequence={} sequence_delta={} cursor_changed={} bracketed_baseline={} bracketed_current={} tail_chars={} probe_chars={} probe_match={} placeholder_baseline={} placeholder_current={} child_live={} foreground={} retained_output={} retained_foreground={} retained_snapshots={} retained_resized={} retained_gaps={} retained_reader_errors={} retained_operator_actions={} retained_exit_code={} output_flags={})",
3963 spec.id,
3964 baseline.sequence,
3965 snapshot.sequence,
3966 snapshot.sequence.saturating_sub(baseline.sequence),
3967 snapshot.cursor != baseline.cursor,
3968 baseline.bracketed_paste,
3969 snapshot.bracketed_paste,
3970 tail.chars().count(),
3971 probe.chars().count(),
3972 !probe.is_empty() && compact_tail.contains(&probe),
3973 paste_placeholder_visible(&baseline_tail),
3974 paste_placeholder_visible(&normalized_tail),
3975 child_live,
3976 foreground_name,
3977 retained.output,
3978 retained.foreground,
3979 retained.snapshots,
3980 retained.resized,
3981 retained.gaps,
3982 retained.reader_errors,
3983 retained.operator_actions,
3984 retained
3985 .exit_code
3986 .map_or_else(|| "none".to_owned(), |code| code.to_string()),
3987 if retained.output_flags.is_empty() {
3988 "none".to_owned()
3989 } else {
3990 retained.output_flags.join(",")
3991 },
3992 ));
3993 }
3994 tokio::time::sleep(Duration::from_millis(STARTUP_GATE_POLL_MS)).await;
3995 }
3996}
3997
3998#[derive(Default)]
3999struct RenderAckEventSummary {
4000 output: usize,
4001 foreground: usize,
4002 snapshots: usize,
4003 resized: usize,
4004 gaps: usize,
4005 reader_errors: usize,
4006 operator_actions: usize,
4007 exit_code: Option<i32>,
4008 output_flags: Vec<&'static str>,
4009}
4010
4011fn render_ack_event_summary(session: &PtySession, baseline_sequence: u64) -> RenderAckEventSummary {
4012 let mut summary = RenderAckEventSummary::default();
4013 let Ok(attachment) = session.attach_retained_events() else {
4014 return summary;
4015 };
4016 for envelope in attachment
4017 .replay
4018 .into_iter()
4019 .filter(|envelope| envelope.sequence > baseline_sequence)
4020 {
4021 match envelope.event {
4022 PtyEvent::Output(bytes) => {
4023 summary.output = summary.output.saturating_add(1);
4024 observe_render_ack_output_flags(&mut summary.output_flags, &bytes);
4025 }
4026 PtyEvent::ForegroundProcess(_) => {
4027 summary.foreground = summary.foreground.saturating_add(1);
4028 }
4029 PtyEvent::SnapshotAvailable { .. } => {
4030 summary.snapshots = summary.snapshots.saturating_add(1);
4031 }
4032 PtyEvent::Resized(_) => {
4033 summary.resized = summary.resized.saturating_add(1);
4034 }
4035 PtyEvent::DataGap { .. } => {
4036 summary.gaps = summary.gaps.saturating_add(1);
4037 }
4038 PtyEvent::ReaderError { .. } => {
4039 summary.reader_errors = summary.reader_errors.saturating_add(1);
4040 }
4041 PtyEvent::OperatorActionRequired { .. } => {
4042 summary.operator_actions = summary.operator_actions.saturating_add(1);
4043 }
4044 PtyEvent::Exited { code } => summary.exit_code = Some(code),
4045 PtyEvent::Started => {}
4046 }
4047 }
4048 summary
4049}
4050
4051fn observe_render_ack_output_flags(flags: &mut Vec<&'static str>, bytes: &[u8]) {
4052 let normalized = String::from_utf8_lossy(bytes).to_ascii_lowercase();
4053 for (label, marker) in [
4054 ("error", "error"),
4055 ("panic", "panic"),
4056 ("fatal", "fatal"),
4057 ("login", "login"),
4058 ("auth", "auth"),
4059 ("permission", "permission"),
4060 ("rate-limit", "rate limit"),
4061 ("usage-limit", "usage limit"),
4062 ("update", "update"),
4063 ("working", "working"),
4064 ("thinking", "thinking"),
4065 ] {
4066 if normalized.contains(marker) && !flags.contains(&label) {
4067 flags.push(label);
4068 }
4069 }
4070}
4071
4072fn safe_process_label(process_name: &str) -> String {
4073 process_name
4074 .chars()
4075 .take(64)
4076 .map(|character| {
4077 if character.is_ascii_alphanumeric() || matches!(character, '.' | '_' | '-') {
4078 character
4079 } else {
4080 '?'
4081 }
4082 })
4083 .collect()
4084}
4085
4086fn prompt_rendered(
4087 snapshot: &PtyTerminalSnapshot,
4088 baseline: &PtyTerminalSnapshot,
4089 probe: &str,
4090) -> bool {
4091 if snapshot.sequence <= baseline.sequence {
4092 return false;
4093 }
4094 let tail = terminal_tail(&snapshot.contents);
4095 let compact = compact_alphanumeric(&tail);
4096 if !probe.is_empty() && compact.contains(probe) {
4097 return true;
4098 }
4099 let normalized_tail = tail.to_ascii_lowercase();
4100 let normalized_baseline = terminal_tail(&baseline.contents).to_ascii_lowercase();
4101 paste_placeholder_visible(&normalized_tail)
4102 && !paste_placeholder_visible(&normalized_baseline)
4103}
4104
4105fn paste_placeholder_visible(normalized_tail: &str) -> bool {
4106 normalized_tail.contains("[pasted content")
4110 || normalized_tail.contains("[pasted text #")
4111}
4112
4113fn terminal_tail(contents: &str) -> String {
4114 let mut lines = contents.lines().rev().take(12).collect::<Vec<_>>();
4115 lines.reverse();
4116 lines.join("\n")
4117}
4118
4119fn prompt_render_probe(prompt: &str) -> String {
4120 let compact = compact_alphanumeric(prompt);
4121 let chars = compact.chars().collect::<Vec<_>>();
4122 chars[chars.len().saturating_sub(32)..].iter().collect()
4123}
4124
4125fn compact_alphanumeric(value: &str) -> String {
4126 value
4127 .chars()
4128 .filter(|character| character.is_alphanumeric())
4129 .flat_map(char::to_lowercase)
4130 .collect()
4131}
4132
4133async fn wait_for_startup_operator_gate(
4134 session: &PtySession,
4135 spec: &AgentSpec,
4136) -> Result<(), String> {
4137 let deadline = Instant::now() + Duration::from_millis(STARTUP_GATE_SETTLE_MS);
4138 loop {
4139 ensure_no_startup_operator_gate(session, spec)?;
4140 let now = Instant::now();
4141 if now >= deadline {
4142 return Ok(());
4143 }
4144 tokio::time::sleep(
4145 deadline
4146 .saturating_duration_since(now)
4147 .min(Duration::from_millis(STARTUP_GATE_POLL_MS)),
4148 )
4149 .await;
4150 }
4151}
4152
4153fn ensure_no_startup_operator_gate(session: &PtySession, spec: &AgentSpec) -> Result<(), String> {
4173 let snapshot = session.terminal_state().map_err(|error| error.to_string())?;
4174 match startup_operator_gate(&snapshot.contents) {
4175 Some(gate) => Err(startup_operator_error(spec, &gate)),
4176 None => Ok(()),
4177 }
4178}
4179
4180fn startup_operator_error(spec: &AgentSpec, gate: &OperatorGateState) -> String {
4181 format!(
4182 "agent '{}' requires operator action at startup ({gate}); initial prompt was not submitted",
4183 spec.id
4184 )
4185}
4186
4187fn ensure_no_readiness_operator_gate(
4188 diagnostics: &ReadinessDiagnostics,
4189 spec: &AgentSpec,
4190) -> Result<(), String> {
4191 match &diagnostics.operator_gate {
4192 Some(gate) => Err(startup_operator_error(spec, gate)),
4193 None => Ok(()),
4194 }
4195}
4196
4197fn normalize_screen_text(contents: &str) -> String {
4204 contents
4205 .split_whitespace()
4206 .collect::<Vec<_>>()
4207 .join(" ")
4208 .to_ascii_lowercase()
4209}
4210
4211const PACKAGE_MANAGER_INSTALL_MARKERS: &[&str] = &[
4217 "npm install -g",
4218 "npm i -g",
4219 "npm update -g",
4220 "pip install --upgrade",
4221 "pip install -u",
4222 "yarn global add",
4223 "pnpm add -g",
4224];
4225
4226const PACKAGE_MANAGER_UPDATE_COMPLETION_MARKERS: &[&str] = &[
4232 "update ran successfully",
4233 "please restart",
4234 "update complete",
4235 "updated successfully",
4236];
4237
4238fn operator_gate(kind: OperatorGateKind, subject: OperatorGateSubject, contents: &str) -> OperatorGateState {
4249 let (input, options) = parse_operator_gate_options(contents);
4250 OperatorGateState::new(kind)
4251 .with_subject(subject)
4252 .with_options(input, options)
4253}
4254
4255fn startup_operator_gate(contents: &str) -> Option<OperatorGateState> {
4256 let normalized = normalize_screen_text(contents);
4257 if [
4258 "trust this folder",
4259 "trust the files in this folder",
4260 "trust the contents of this directory",
4261 "do you trust this directory",
4262 ]
4263 .iter()
4264 .any(|marker| normalized.contains(marker))
4265 {
4266 return Some(operator_gate(
4267 OperatorGateKind::WorkspaceTrust,
4268 OperatorGateSubject::Directory { path: None },
4269 contents,
4270 ));
4271 }
4272 if normalized.contains("quick safety check")
4273 && (normalized.contains("yes, i trust this folder")
4274 || normalized.contains("continue without these permissions"))
4275 {
4276 return Some(operator_gate(
4277 OperatorGateKind::WorkspaceTrust,
4278 OperatorGateSubject::Directory { path: None },
4279 contents,
4280 ));
4281 }
4282 if normalized.contains("hooks need review") && normalized.contains("continue without trusting")
4297 {
4298 return Some(operator_gate(
4299 OperatorGateKind::HookTrust,
4300 OperatorGateSubject::Hooks { count: None },
4301 contents,
4302 ));
4303 }
4304 if [
4305 "select authentication method",
4306 "choose how to authenticate",
4307 "no auth type is selected",
4308 ]
4309 .iter()
4310 .any(|marker| normalized.contains(marker))
4311 {
4312 return Some(operator_gate(
4313 OperatorGateKind::Authentication,
4314 OperatorGateSubject::Account,
4315 contents,
4316 ));
4317 }
4318 if normalized.contains("sign in")
4319 && (normalized.contains("openai")
4320 || normalized.contains("chatgpt")
4321 || normalized.contains("codex"))
4322 {
4323 return Some(operator_gate(
4324 OperatorGateKind::Authentication,
4325 OperatorGateSubject::Account,
4326 contents,
4327 ));
4328 }
4329 if normalized.contains("select login method") {
4338 return Some(operator_gate(
4339 OperatorGateKind::Authentication,
4340 OperatorGateSubject::Account,
4341 contents,
4342 ));
4343 }
4344 if normalized.contains("sign in") && normalized.contains("paste code") {
4356 return Some(
4357 OperatorGateState::new(OperatorGateKind::Authentication)
4358 .with_subject(OperatorGateSubject::Account)
4359 .with_options(OperatorGateInput::TextEntry, Vec::new()),
4360 );
4361 }
4362 if normalized.contains("kimi code update available")
4363 && normalized.contains("install update now")
4364 {
4365 return Some(operator_gate(
4366 OperatorGateKind::VendorUpdate,
4367 OperatorGateSubject::Unknown,
4368 contents,
4369 ));
4370 }
4371 if PACKAGE_MANAGER_INSTALL_MARKERS
4382 .iter()
4383 .any(|marker| normalized.contains(marker))
4384 && PACKAGE_MANAGER_UPDATE_COMPLETION_MARKERS
4385 .iter()
4386 .any(|marker| normalized.contains(marker))
4387 {
4388 return Some(operator_gate(
4389 OperatorGateKind::VendorUpdate,
4390 OperatorGateSubject::Unknown,
4391 contents,
4392 ));
4393 }
4394 if normalized.contains("choose the text style that looks best with your terminal") {
4395 return Some(operator_gate(
4396 OperatorGateKind::TerminalAppearance,
4397 OperatorGateSubject::Appearance,
4398 contents,
4399 ));
4400 }
4401 if normalized.contains("welcome to claude code for")
4402 && (normalized.contains("open files") || normalized.contains("selected lines"))
4403 {
4404 return Some(operator_gate(
4405 OperatorGateKind::Onboarding,
4406 OperatorGateSubject::Unknown,
4407 contents,
4408 ));
4409 }
4410 if normalized.contains("welcome to claude code")
4411 && (normalized.contains("press enter") || normalized.contains("enter to continue"))
4412 {
4413 return Some(operator_gate(
4414 OperatorGateKind::Onboarding,
4415 OperatorGateSubject::Unknown,
4416 contents,
4417 ));
4418 }
4419 if normalized.contains("migration") && normalized.contains("enter confirm") && normalized.contains("esc")
4420 {
4421 return Some(operator_gate(
4422 OperatorGateKind::ConfigurationMigration,
4423 OperatorGateSubject::Unknown,
4424 contents,
4425 ));
4426 }
4427 None
4428}
4429
4430fn parse_operator_gate_options(contents: &str) -> (OperatorGateInput, Vec<OperatorGateOption>) {
4452 let numbered: Vec<OperatorGateOption> = contents
4453 .lines()
4454 .filter_map(parse_numbered_gate_option_line)
4455 .take(OPERATOR_GATE_OPTIONS_MAX)
4456 .collect();
4457 if !numbered.is_empty() {
4458 return (OperatorGateInput::NumberedList, numbered);
4459 }
4460 let arrow: Vec<OperatorGateOption> = contents
4461 .lines()
4462 .filter_map(parse_arrow_gate_option_line)
4463 .take(OPERATOR_GATE_OPTIONS_MAX)
4464 .collect();
4465 if !arrow.is_empty() {
4466 return (OperatorGateInput::ArrowList, arrow);
4467 }
4468 (OperatorGateInput::Unknown, Vec::new())
4469}
4470
4471const GATE_CURSOR_GLYPHS: [char; 2] = ['\u{203a}', '\u{276f}'];
4485
4486fn parse_numbered_gate_option_line(line: &str) -> Option<OperatorGateOption> {
4487 let trimmed = line.trim_start();
4488 let (selected, rest) = match trimmed.strip_prefix(GATE_CURSOR_GLYPHS) {
4496 Some(stripped) => (true, stripped.trim_start()),
4497 None => (false, trimmed),
4498 };
4499 let digits_end = rest.find(|character: char| !character.is_ascii_digit()).unwrap_or(0);
4500 if digits_end == 0 {
4501 return None;
4502 }
4503 let text = rest[digits_end..].strip_prefix('.')?.trim();
4504 if text.is_empty() {
4505 return None;
4506 }
4507 Some(OperatorGateOption {
4508 semantics: classify_operator_gate_option_semantics(text),
4509 text: text.to_owned(),
4510 selected,
4511 })
4512}
4513
4514fn parse_arrow_gate_option_line(line: &str) -> Option<OperatorGateOption> {
4538 let trimmed = line.trim();
4539 if trimmed.is_empty()
4540 || trimmed
4541 .chars()
4542 .any(|character| matches!(character, '.' | '?' | '!' | ':'))
4543 {
4544 return None;
4545 }
4546 let lower = trimmed.to_ascii_lowercase();
4547 if lower.contains('\u{2191}') || lower.contains('\u{2193}') || lower.contains("navigate") {
4548 return None;
4549 }
4550 let (selected, rest) = match trimmed.strip_prefix(GATE_CURSOR_GLYPHS) {
4551 Some(stripped) => (true, stripped.trim()),
4552 None => (false, trimmed),
4553 };
4554 if rest.is_empty() || !rest.chars().any(|character| character.is_alphabetic()) {
4555 return None;
4556 }
4557 let semantics = classify_operator_gate_option_semantics(rest);
4558 if !selected && semantics == OperatorGateOptionSemantics::Unknown {
4559 return None;
4560 }
4561 Some(OperatorGateOption {
4562 semantics,
4563 text: rest.to_owned(),
4564 selected,
4565 })
4566}
4567
4568fn classify_operator_gate_option_semantics(text: &str) -> OperatorGateOptionSemantics {
4576 let lower = text.to_ascii_lowercase();
4577 if lower.contains("review") {
4578 OperatorGateOptionSemantics::Inspect
4579 } else if lower.contains("don't trust")
4580 || lower.contains("do not trust")
4581 || lower.contains("without trusting")
4582 || lower.contains("decline")
4583 {
4584 OperatorGateOptionSemantics::Decline
4585 } else if lower.contains("trust") || lower.contains("continue") || lower.contains("proceed") {
4586 OperatorGateOptionSemantics::Accept
4587 } else if lower.contains("exit") || lower.contains("quit") || lower.contains("cancel") {
4588 OperatorGateOptionSemantics::Exit
4589 } else {
4590 OperatorGateOptionSemantics::Unknown
4591 }
4592}
4593
4594fn looks_like_shell_prompt_line(line: &str) -> bool {
4599 matches!(line.trim_end().chars().last(), Some('$' | '%' | '#' | '>'))
4600}
4601
4602fn screen_failure(contents: &str) -> Option<&'static str> {
4625 let normalized = normalize_screen_text(contents);
4626
4627 if ["panicked at", "unhandled exception", "traceback (most recent call last)"]
4636 .iter()
4637 .any(|marker| normalized.contains(marker))
4638 {
4639 return Some("crash");
4640 }
4641
4642 if normalized.contains("command not found")
4648 || normalized.contains("is not recognized as an internal or external command")
4649 || (normalized.contains("no such file or directory")
4650 && contents.lines().any(looks_like_shell_prompt_line))
4651 {
4652 return Some("missing command");
4653 }
4654
4655 if normalized.contains("api key required") && normalized.contains("grok_api_key") {
4666 return Some("missing api key");
4667 }
4668
4669 None
4670}
4671
4672fn screen_failure_for_generation(contents: &str, ever_reached_ready: bool) -> Option<&'static str> {
4679 if ever_reached_ready {
4680 None
4681 } else {
4682 screen_failure(contents)
4683 }
4684}
4685
4686#[derive(Clone, Debug, Eq, PartialEq)]
4689enum ForegroundVerdict {
4690 Agent,
4692 Foreign { process: String },
4694}
4695
4696fn classify_pty_screen_state(
4722 foreground: Option<&ForegroundVerdict>,
4723 gate: Option<&OperatorGateState>,
4724 failure: Option<&'static str>,
4725) -> PtyScreenState {
4726 if let Some(ForegroundVerdict::Foreign { process }) = foreground {
4727 return PtyScreenState::NotAgent {
4728 observed_process: process.clone(),
4729 };
4730 }
4731 if let Some(reason) = failure {
4732 return PtyScreenState::Failing {
4733 reason: reason.to_owned(),
4734 };
4735 }
4736 if let Some(gate) = gate {
4737 return PtyScreenState::OperatorGate { gate: gate.clone() };
4738 }
4739 if foreground.is_some() {
4740 PtyScreenState::Ready
4741 } else {
4742 PtyScreenState::Unknown
4743 }
4744}
4745
4746fn resolve_foreground_verdict(
4762 spec: &AgentSpec,
4763 observation: &PtyForegroundObservation,
4764 platform: RuntimePlatform,
4765) -> ForegroundVerdict {
4766 let process_name = observation
4767 .readiness
4768 .process_name
4769 .as_deref()
4770 .unwrap_or(observation.observed_process.as_str());
4771 if is_expected_agent_process(spec, process_name, platform)
4772 || is_agent_foreground_wrapper(process_name, platform)
4773 {
4774 ForegroundVerdict::Agent
4775 } else {
4776 ForegroundVerdict::Foreign {
4777 process: observation.observed_process.clone(),
4778 }
4779 }
4780}
4781
4782#[derive(Clone, Copy, Debug, Eq, PartialEq)]
4788enum ForegroundProbeSchedule {
4789 Disarmed,
4791 Armed,
4793}
4794
4795fn foreground_probe_schedule(state: &PtyScreenState) -> ForegroundProbeSchedule {
4796 if matches!(state, PtyScreenState::Ready) {
4797 ForegroundProbeSchedule::Disarmed
4798 } else {
4799 ForegroundProbeSchedule::Armed
4800 }
4801}
4802
4803fn foreground_probe_rearms_immediately(previous: &PtyScreenState, next: &PtyScreenState) -> bool {
4811 matches!(previous, PtyScreenState::Ready) && !matches!(next, PtyScreenState::Ready)
4812}
4813
4814fn observe_readiness_event(
4815 tracker: &mut ReadinessTracker<'_>,
4816 diagnostics: &mut ReadinessDiagnostics,
4817 event: PtyEvent,
4818 elapsed_ms: u64,
4819) -> Result<(), String> {
4820 match event {
4821 PtyEvent::Output(data) => {
4822 diagnostics.observe_output(&data);
4823 tracker.observe_output(&data, elapsed_ms);
4824 }
4825 PtyEvent::ForegroundProcess(observation) => {
4826 diagnostics.observe_foreground(&observation.readiness);
4827 tracker.observe_foreground(&observation.readiness, elapsed_ms);
4828 }
4829 PtyEvent::DataGap { .. } => {
4830 return Err("PTY replay gap prevents positive readiness proof".to_owned());
4831 }
4832 PtyEvent::ReaderError { message } | PtyEvent::OperatorActionRequired { message } => {
4833 return Err(message);
4834 }
4835 PtyEvent::Exited { code } => {
4836 return Err(format!("PTY exited with code {code} before readiness"));
4837 }
4838 PtyEvent::Started | PtyEvent::Resized(_) | PtyEvent::SnapshotAvailable { .. } => {}
4839 }
4840 Ok(())
4841}
4842
4843#[derive(Default)]
4844struct ReadinessDiagnostics {
4845 draft_signal: Option<gate4agent_types::DraftReadySignal>,
4846 output_bytes: usize,
4847 output_chunks: usize,
4848 tail: Vec<u8>,
4849 saw_bracketed_paste: bool,
4850 saw_cursor_show: bool,
4851 saw_cursor_hide: bool,
4852 saw_alternate_screen: bool,
4853 saw_clear_screen: bool,
4854 saw_claude_composer: bool,
4855 saw_codex_composer: bool,
4856 saw_named_foreground: bool,
4857 operator_gate: Option<OperatorGateState>,
4858}
4859
4860impl ReadinessDiagnostics {
4861 fn observe_output(&mut self, data: &[u8]) {
4862 const SIGNAL_TAIL_BYTES: usize = 4_096;
4863 self.output_bytes = self.output_bytes.saturating_add(data.len());
4864 self.output_chunks = self.output_chunks.saturating_add(1);
4865 let mut combined = Vec::with_capacity(self.tail.len().saturating_add(data.len()));
4866 combined.extend_from_slice(&self.tail);
4867 combined.extend_from_slice(data);
4868 self.saw_bracketed_paste |= readiness_bytes_contain(&combined, b"\x1b[?2004h");
4869 self.saw_cursor_show |= readiness_bytes_contain(&combined, b"\x1b[?25h");
4870 self.saw_cursor_hide |= readiness_bytes_contain(&combined, b"\x1b[?25l");
4871 self.saw_alternate_screen |= readiness_bytes_contain(&combined, b"\x1b[?1049h");
4872 self.saw_clear_screen |= readiness_bytes_contain(&combined, b"\x1b[2J");
4873 self.saw_claude_composer |= readiness_bytes_contain(&combined, "❯".as_bytes());
4874 self.saw_codex_composer |= readiness_bytes_contain(&combined, "›".as_bytes());
4875 let text = String::from_utf8_lossy(&combined);
4876 self.operator_gate = self
4877 .operator_gate
4878 .take()
4879 .or_else(|| startup_operator_gate(&strip_ansi_codes(&text)));
4880 self.tail = combined[combined.len().saturating_sub(SIGNAL_TAIL_BYTES)..].to_vec();
4881 }
4882
4883 fn observe_foreground(&mut self, foreground: &gate4agent::agent::ForegroundObservation) {
4884 self.saw_named_foreground |= foreground.process_name.is_some();
4885 }
4886
4887 fn summary(&self) -> String {
4888 format!(
4889 "draft_signal={:?} output_bytes={} output_chunks={} named_foreground={} bracketed_paste={} cursor_show={} cursor_hide={} alternate_screen={} clear_screen={} claude_composer={} codex_composer={} csi={}",
4890 self.draft_signal,
4891 self.output_bytes,
4892 self.output_chunks,
4893 self.saw_named_foreground,
4894 self.saw_bracketed_paste,
4895 self.saw_cursor_show,
4896 self.saw_cursor_hide,
4897 self.saw_alternate_screen,
4898 self.saw_clear_screen,
4899 self.saw_claude_composer,
4900 self.saw_codex_composer,
4901 readiness_csi_signatures(&self.tail),
4902 )
4903 }
4904}
4905
4906fn readiness_csi_signatures(bytes: &[u8]) -> String {
4907 let mut signatures = Vec::new();
4908 let mut index = 0;
4909 while index + 2 < bytes.len() && signatures.len() < 32 {
4910 if bytes[index] != 0x1b || bytes[index + 1] != b'[' {
4911 index += 1;
4912 continue;
4913 }
4914 let mut end = index + 2;
4915 while end < bytes.len() && end.saturating_sub(index) <= 24 {
4916 let byte = bytes[end];
4917 if (0x40..=0x7e).contains(&byte) {
4918 let signature = String::from_utf8_lossy(&bytes[index + 2..=end]).into_owned();
4919 if !signatures.iter().any(|existing| existing == &signature) {
4920 signatures.push(signature);
4921 }
4922 index = end;
4923 break;
4924 }
4925 if !(0x20..=0x3f).contains(&byte) {
4926 break;
4927 }
4928 end += 1;
4929 }
4930 index += 1;
4931 }
4932 if signatures.is_empty() {
4933 "none".to_owned()
4934 } else {
4935 signatures.join("|")
4936 }
4937}
4938
4939fn readiness_bytes_contain(haystack: &[u8], needle: &[u8]) -> bool {
4940 haystack
4941 .windows(needle.len())
4942 .any(|window| window == needle)
4943}
4944
4945fn readiness_complete(
4946 status: ReadinessStatus,
4947 diagnostics: &ReadinessDiagnostics,
4948) -> Result<bool, String> {
4949 match status {
4950 ReadinessStatus::Waiting => Ok(false),
4951 ReadinessStatus::Ready(_) => Ok(true),
4952 ReadinessStatus::TimedOut => Err(format!(
4953 "PTY readiness timed out ({})",
4954 diagnostics.summary()
4955 )),
4956 }
4957}
4958
4959fn elapsed_ms(started: Instant) -> u64 {
4960 u64::try_from(started.elapsed().as_millis()).unwrap_or(u64::MAX)
4961}
4962
4963#[cfg(test)]
4964mod tests {
4965 use super::{
4966 acp_approval_level_args,
4967 approval_level_not_offered_message,
4968 classify_operator_gate_option_semantics,
4969 classify_pty_screen_state,
4970 defers_permission_requests,
4971 foreground_probe_rearms_immediately, foreground_probe_schedule,
4972 argument_looks_like_credential, host_policy_for_approval_level,
4973 mcp_server_acp_entry,
4974 parse_operator_gate_options, prepare_fresh_pty_provider_session,
4975 prompt_render_probe, prompt_rendered, required_acp_mode,
4976 reserve_provider_gap_sequence, resolve_foreground_verdict, screen_failure,
4977 screen_failure_for_generation, with_pty_terminal_capability_defaults,
4978 should_attach_pty_provider_stream, should_probe_pty_identity, startup_operator_gate,
4979 terminal_frame, terminal_frame_byte_len, terminal_state_capture_should_skip,
4980 validate_instance_launch_arguments, validate_spawn_runtime_policy,
4981 ForegroundProbeSchedule, ForegroundVerdict, ReadinessDiagnostics, RateLimitFeed,
4982 Utf8ChunkDecoder,
4983 };
4984 use gate4agent::acp::protocol::{McpServerConfig, SessionMode};
4985 use gate4agent::HostPolicy;
4986 use gate4agent_adapters::builtin_adapter_registry;
4987 use gate4agent_catalog::{EnvMutation, McpServerSpec, ModeId};
4988 use gate4agent::agent::ForegroundObservation;
4989 use gate4agent::core::types::{
4990 AgentEvent, ContextWindowUsage as AgentContextWindowUsage, HostDecisionAuthority,
4991 HostRequestDecision, HostRequestOutcome, RateLimitType,
4992 };
4993 use gate4agent::pty::event::PtyMouseProtocolEncoding;
4994 use gate4agent::pty::{PtyForegroundObservation, PtyForegroundSource, RateLimitDetector};
4995 use gate4agent::CliTool;
4996 use gate4agent_types::{
4997 AdapterFamily, AgentId, ApprovalLevel,
4998 HostDecisionAuthority as ProviderHostDecisionAuthority,
4999 HostRequestDecision as ProviderHostRequestDecision,
5000 HostRequestOutcome as ProviderHostRequestOutcome, OperatorGateInput,
5001 OperatorGateKind, OperatorGateOptionSemantics, OperatorGateState, OperatorGateSubject,
5002 ProviderEvent, ProviderInteractionKind, ProviderInteractionOption, ProviderRuntimePolicy,
5003 PtyScreenState, RuntimePlatform, TerminalMouseProtocolEncoding, TransportKind,
5004 };
5005 use std::ffi::{OsStr, OsString};
5006
5007 fn snapshot(sequence: u64, contents: &str) -> super::PtyTerminalSnapshot {
5008 super::PtyTerminalSnapshot {
5009 pty_id: "fixture".to_owned(),
5010 provider_revision: "fixture-r1".to_owned(),
5011 generation: 1,
5012 sequence,
5013 size: gate4agent::pty::PtySize { rows: 24, cols: 80 },
5014 cursor: (0, 0),
5015 bracketed_paste: false,
5016 contents: contents.to_owned(),
5017 formatted: Vec::new(),
5018 scrollback_formatted: Vec::new(),
5019 alternate_screen: false,
5020 mouse_protocol_enabled: false,
5021 mouse_protocol_encoding: PtyMouseProtocolEncoding::Default,
5022 produced_at_unix_ms: 0,
5023 }
5024 }
5025
5026 #[test]
5027 fn terminal_frame_preserves_scrollback_and_terminal_input_metadata() {
5028 let mut snapshot = snapshot(7, "visible");
5029 snapshot.scrollback_formatted = vec![b"older".to_vec()];
5030 snapshot.alternate_screen = true;
5031 snapshot.mouse_protocol_enabled = true;
5032 snapshot.mouse_protocol_encoding = PtyMouseProtocolEncoding::Sgr;
5033 snapshot.produced_at_unix_ms = 1_700_000_000_000;
5034
5035 let frame = terminal_frame(snapshot, PtyScreenState::default());
5036 assert_eq!(frame.scrollback_formatted, vec![b"older".to_vec()]);
5037 assert!(frame.alternate_screen);
5038 assert!(frame.mouse_protocol_enabled);
5039 assert_eq!(frame.mouse_protocol_encoding, TerminalMouseProtocolEncoding::Sgr);
5040 assert_eq!(frame.produced_at_unix_ms, 1_700_000_000_000);
5042 }
5043
5044 #[test]
5045 fn terminal_frame_byte_len_sums_formatted_and_every_scrollback_row() {
5046 let mut snapshot = snapshot(9, "visible");
5047 snapshot.formatted = b"\x1b[2Jvisible".to_vec();
5048 snapshot.scrollback_formatted =
5049 vec![b"row one".to_vec(), b"row two, longer".to_vec()];
5050 let expected = snapshot.formatted.len()
5051 + snapshot.scrollback_formatted[0].len()
5052 + snapshot.scrollback_formatted[1].len();
5053
5054 let frame = terminal_frame(snapshot, PtyScreenState::default());
5055 assert_eq!(terminal_frame_byte_len(&frame), expected as u64);
5056 }
5057
5058 #[test]
5059 fn the_sequence_gate_skips_only_when_the_sequence_did_not_advance() {
5060 assert!(terminal_state_capture_should_skip(Ok::<u64, ()>(5), 5));
5061 assert!(terminal_state_capture_should_skip(Ok::<u64, ()>(4), 5));
5062 assert!(!terminal_state_capture_should_skip(Ok::<u64, ()>(6), 5));
5063 assert!(!terminal_state_capture_should_skip(Err::<u64, ()>(()), 5));
5066 }
5067
5068 #[test]
5069 fn structured_context_usage_maps_without_pty_inference() {
5070 let mapped = super::provider_event(
5071 AgentEvent::ContextWindowUsage {
5072 usage: AgentContextWindowUsage {
5073 uncached_input_tokens: 70,
5074 cache_read_tokens: 20,
5075 cache_write_tokens: 0,
5076 output_tokens: 10,
5077 unattributed_tokens: 5,
5078 used_tokens: 105,
5079 capacity_tokens: 100,
5080 },
5081 },
5082 &[],
5083 );
5084 assert_eq!(
5085 mapped,
5086 Some(ProviderEvent::ContextWindowUsage {
5087 usage: gate4agent_types::ContextWindowUsage {
5088 uncached_input_tokens: 70,
5089 cache_read_tokens: 20,
5090 cache_write_tokens: 0,
5091 output_tokens: 10,
5092 unattributed_tokens: 5,
5093 used_tokens: 105,
5094 capacity_tokens: 100,
5095 }
5096 })
5097 );
5098 assert_eq!(
5099 ProviderEvent::ContextWindowUsage {
5100 usage: gate4agent_types::ContextWindowUsage {
5101 uncached_input_tokens: 70,
5102 cache_read_tokens: 20,
5103 cache_write_tokens: 0,
5104 output_tokens: 10,
5105 unattributed_tokens: 5,
5106 used_tokens: 104,
5107 capacity_tokens: 100,
5108 },
5109 }
5110 .validate_ingress(),
5111 Err(gate4agent_types::ProviderEventValidationError::ContextWindowSegmentsMismatch {
5112 segment_sum: 105,
5113 used_tokens: 104,
5114 })
5115 );
5116 assert!(
5117 super::provider_event(AgentEvent::PtyRaw { data: b"105/100".to_vec() }, &[]).is_none()
5118 );
5119 }
5120
5121 #[test]
5130 fn turn_interrupted_maps_to_the_provider_event_that_unblocks_the_next_prompt() {
5131 let mapped = super::provider_event(
5132 AgentEvent::TurnInterrupted {
5133 reason: "session/prompt timed out after 120s".to_owned(),
5134 },
5135 &[],
5136 );
5137 assert_eq!(mapped, Some(ProviderEvent::TurnInterrupted));
5138 }
5139
5140 #[test]
5141 fn tool_result_non_execution_kind_carries_through_to_tool_completed() {
5142 let mapped = super::provider_event(
5143 AgentEvent::ToolResult {
5144 id: "t1".to_owned(),
5145 output: "denied".to_owned(),
5146 is_error: true,
5147 duration_ms: None,
5148 non_execution_kind: Some("permission-rule".to_owned()),
5149 },
5150 &[],
5151 );
5152 assert_eq!(
5153 mapped,
5154 Some(ProviderEvent::ToolCompleted {
5155 id: "t1".to_owned(),
5156 output: "denied".to_owned(),
5157 is_error: true,
5158 duration_ms: None,
5159 agent_id: None,
5160 non_execution_kind: Some("permission-rule".to_owned()),
5161 })
5162 );
5163 }
5164
5165 #[test]
5166 fn session_end_stop_reason_carries_through_to_session_ended() {
5167 use gate4agent::StopReason;
5168
5169 let refusal = super::provider_event(
5170 AgentEvent::SessionEnd {
5171 result: "refusal".to_owned(),
5172 cost_usd: None,
5173 is_error: true,
5174 stop_reason: Some(StopReason::Refusal),
5175 },
5176 &[],
5177 );
5178 assert_eq!(
5179 refusal,
5180 Some(ProviderEvent::SessionEnded {
5181 result: "refusal".to_owned(),
5182 cost_usd: None,
5183 is_error: true,
5184 stop_reason: Some(gate4agent_types::ProviderStopReason::Refusal),
5185 })
5186 );
5187
5188 let quota = super::provider_event(
5189 AgentEvent::SessionEnd {
5190 result: "JSON-RPC error -32603: Internal error".to_owned(),
5191 cost_usd: None,
5192 is_error: true,
5193 stop_reason: Some(StopReason::ProviderError {
5194 code: -32603,
5195 message: "You've hit your usage limit.".to_owned(),
5196 vendor_code: Some("usageLimitExceeded".to_owned()),
5197 }),
5198 },
5199 &[],
5200 );
5201 assert_eq!(
5202 quota,
5203 Some(ProviderEvent::SessionEnded {
5204 result: "JSON-RPC error -32603: Internal error".to_owned(),
5205 cost_usd: None,
5206 is_error: true,
5207 stop_reason: Some(gate4agent_types::ProviderStopReason::ProviderError {
5208 code: -32603,
5209 message: "You've hit your usage limit.".to_owned(),
5210 vendor_code: Some("usageLimitExceeded".to_owned()),
5211 }),
5212 })
5213 );
5214 }
5215
5216 #[test]
5217 fn host_request_reaches_the_operator_with_method_and_host_decision() {
5218 let denied = super::provider_event(
5219 AgentEvent::RpcIncomingRequest {
5220 id: gate4agent::rpc::message::RpcId::Number(1),
5221 method: "fs/read_text_file".to_owned(),
5222 params: None,
5223 decision: HostRequestDecision::Denied { by: HostDecisionAuthority::Policy },
5224 outcome: HostRequestOutcome::Executed,
5225 reason: None,
5226 },
5227 &[],
5228 );
5229 assert_eq!(
5230 denied,
5231 Some(ProviderEvent::HostRequestObserved {
5232 method: "fs/read_text_file".to_owned(),
5233 params_json: String::new(),
5234 decision: ProviderHostRequestDecision::Denied {
5235 by: ProviderHostDecisionAuthority::Policy
5236 },
5237 outcome: ProviderHostRequestOutcome::Executed,
5238 reason: None,
5239 })
5240 );
5241
5242 let granted = super::provider_event(
5243 AgentEvent::RpcIncomingRequest {
5244 id: gate4agent::rpc::message::RpcId::Number(2),
5245 method: "terminal/create".to_owned(),
5246 params: None,
5247 decision: HostRequestDecision::Granted { by: HostDecisionAuthority::Policy },
5248 outcome: HostRequestOutcome::Executed,
5249 reason: None,
5250 },
5251 &[],
5252 );
5253 assert_eq!(
5254 granted,
5255 Some(ProviderEvent::HostRequestObserved {
5256 method: "terminal/create".to_owned(),
5257 params_json: String::new(),
5258 decision: ProviderHostRequestDecision::Granted {
5259 by: ProviderHostDecisionAuthority::Policy
5260 },
5261 outcome: ProviderHostRequestOutcome::Executed,
5262 reason: None,
5263 })
5264 );
5265
5266 let gate_denied = super::provider_event(
5270 AgentEvent::RpcIncomingRequest {
5271 id: gate4agent::rpc::message::RpcId::Number(3),
5272 method: "terminal/create".to_owned(),
5273 params: None,
5274 decision: HostRequestDecision::Denied { by: HostDecisionAuthority::Gate },
5275 outcome: HostRequestOutcome::Executed,
5276 reason: None,
5277 },
5278 &[],
5279 );
5280 assert_eq!(
5281 gate_denied,
5282 Some(ProviderEvent::HostRequestObserved {
5283 method: "terminal/create".to_owned(),
5284 params_json: String::new(),
5285 decision: ProviderHostRequestDecision::Denied {
5286 by: ProviderHostDecisionAuthority::Gate
5287 },
5288 outcome: ProviderHostRequestOutcome::Executed,
5289 reason: None,
5290 })
5291 );
5292
5293 let operator_granted = super::provider_event(
5294 AgentEvent::RpcIncomingRequest {
5295 id: gate4agent::rpc::message::RpcId::Number(4),
5296 method: "session/request_permission".to_owned(),
5297 params: None,
5298 decision: HostRequestDecision::Granted { by: HostDecisionAuthority::Operator },
5299 outcome: HostRequestOutcome::Executed,
5300 reason: None,
5301 },
5302 &[],
5303 );
5304 assert_eq!(
5305 operator_granted,
5306 Some(ProviderEvent::HostRequestObserved {
5307 method: "session/request_permission".to_owned(),
5308 params_json: String::new(),
5309 decision: ProviderHostRequestDecision::Granted {
5310 by: ProviderHostDecisionAuthority::Operator
5311 },
5312 outcome: ProviderHostRequestOutcome::Executed,
5313 reason: None,
5314 })
5315 );
5316
5317 let deadline_denied = super::provider_event(
5318 AgentEvent::RpcIncomingRequest {
5319 id: gate4agent::rpc::message::RpcId::Number(5),
5320 method: "session/request_permission".to_owned(),
5321 params: None,
5322 decision: HostRequestDecision::Denied { by: HostDecisionAuthority::DeadlinePolicy },
5323 outcome: HostRequestOutcome::Executed,
5324 reason: None,
5325 },
5326 &[],
5327 );
5328 assert_eq!(
5329 deadline_denied,
5330 Some(ProviderEvent::HostRequestObserved {
5331 method: "session/request_permission".to_owned(),
5332 params_json: String::new(),
5333 decision: ProviderHostRequestDecision::Denied {
5334 by: ProviderHostDecisionAuthority::DeadlinePolicy
5335 },
5336 outcome: ProviderHostRequestOutcome::Executed,
5337 reason: None,
5338 })
5339 );
5340
5341 }
5347
5348 #[test]
5355 fn granted_host_request_execution_failure_carries_through_as_a_failed_outcome() {
5356 let spawn_error = "terminal/create spawn failed: os error 3";
5357 let mapped = super::provider_event(
5358 AgentEvent::RpcIncomingRequest {
5359 id: gate4agent::rpc::message::RpcId::Number(10),
5360 method: "terminal/create".to_owned(),
5361 params: None,
5362 decision: HostRequestDecision::Granted { by: HostDecisionAuthority::Policy },
5363 outcome: HostRequestOutcome::Failed { error: spawn_error.to_owned() },
5364 reason: None,
5365 },
5366 &[],
5367 );
5368 assert_eq!(
5369 mapped,
5370 Some(ProviderEvent::HostRequestObserved {
5371 method: "terminal/create".to_owned(),
5372 params_json: String::new(),
5373 decision: ProviderHostRequestDecision::Granted {
5374 by: ProviderHostDecisionAuthority::Policy
5375 },
5376 outcome: ProviderHostRequestOutcome::Failed { error: spawn_error.to_owned() },
5377 reason: None,
5378 })
5379 );
5380 mapped.expect("mapped").validate_ingress().expect("bounded free-text error validates");
5381 }
5382
5383 #[test]
5390 fn host_request_reason_carries_through_to_the_operator() {
5391 let gate_text =
5392 "blocked by dangerous-command gate: rule=filesystem-wipe, argument=rm -rf /";
5393 let denied = super::provider_event(
5394 AgentEvent::RpcIncomingRequest {
5395 id: gate4agent::rpc::message::RpcId::Number(9),
5396 method: "terminal/create".to_owned(),
5397 params: None,
5398 decision: HostRequestDecision::Denied { by: HostDecisionAuthority::Gate },
5399 outcome: HostRequestOutcome::Executed,
5400 reason: Some(gate_text.to_owned()),
5401 },
5402 &[],
5403 );
5404 assert_eq!(
5405 denied,
5406 Some(ProviderEvent::HostRequestObserved {
5407 method: "terminal/create".to_owned(),
5408 params_json: String::new(),
5409 decision: ProviderHostRequestDecision::Denied {
5410 by: ProviderHostDecisionAuthority::Gate
5411 },
5412 outcome: ProviderHostRequestOutcome::Executed,
5413 reason: Some(gate_text.to_owned()),
5414 })
5415 );
5416
5417 denied.expect("mapped").validate_ingress().expect("bounded free-text reason validates");
5422 }
5423
5424 #[test]
5425 fn deferred_acp_permission_request_becomes_an_interaction_with_its_id() {
5426 let with_title = super::provider_event(
5427 AgentEvent::RpcIncomingRequest {
5428 id: gate4agent::rpc::message::RpcId::Number(6),
5429 method: "session/request_permission".to_owned(),
5430 params: Some(serde_json::json!({
5431 "sessionId": "s1",
5432 "toolCall": {"toolCallId": "t1", "kind": "edit", "title": "Edit src/main.rs"},
5433 "options": [
5434 {"optionId": "allow", "name": "Allow", "kind": "allow_once"},
5435 {"optionId": "reject", "name": "Reject", "kind": "reject_once"},
5436 ],
5437 })),
5438 decision: HostRequestDecision::Deferred,
5439 outcome: HostRequestOutcome::Executed,
5440 reason: None,
5441 },
5442 &[],
5443 );
5444 assert_eq!(
5453 with_title,
5454 Some(ProviderEvent::InteractionRequested {
5455 request_id: Some("number:6".to_owned()),
5456 interaction_kind: ProviderInteractionKind::Approval,
5457 tool_name: "edit".to_owned(),
5458 title: Some("Edit src/main.rs".to_owned()),
5459 prompt: "Edit src/main.rs".to_owned(),
5460 options: vec![
5461 ProviderInteractionOption {
5462 option_id: "allow".to_owned(),
5463 name: "Allow".to_owned(),
5464 kind: "allow_once".to_owned(),
5465 },
5466 ProviderInteractionOption {
5467 option_id: "reject".to_owned(),
5468 name: "Reject".to_owned(),
5469 kind: "reject_once".to_owned(),
5470 },
5471 ],
5472 agent_id: None,
5473 })
5474 );
5475
5476 let without_title = super::provider_event(
5482 AgentEvent::RpcIncomingRequest {
5483 id: gate4agent::rpc::message::RpcId::String("agent-7".to_owned()),
5484 method: "session/request_permission".to_owned(),
5485 params: None,
5486 decision: HostRequestDecision::Deferred,
5487 outcome: HostRequestOutcome::Executed,
5488 reason: None,
5489 },
5490 &[],
5491 );
5492 assert_eq!(
5493 without_title,
5494 Some(ProviderEvent::InteractionRequested {
5495 request_id: Some("string:agent-7".to_owned()),
5496 interaction_kind: ProviderInteractionKind::Approval,
5497 tool_name: "session/request_permission".to_owned(),
5498 title: None,
5499 prompt: String::new(),
5500 options: Vec::new(),
5501 agent_id: None,
5502 })
5503 );
5504
5505 let resolved = super::provider_event(
5509 AgentEvent::RpcIncomingRequest {
5510 id: gate4agent::rpc::message::RpcId::Number(6),
5511 method: "session/request_permission".to_owned(),
5512 params: None,
5513 decision: HostRequestDecision::Denied { by: HostDecisionAuthority::Operator },
5514 outcome: HostRequestOutcome::Executed,
5515 reason: None,
5516 },
5517 &[],
5518 );
5519 assert_eq!(
5520 resolved,
5521 Some(ProviderEvent::HostRequestObserved {
5522 method: "session/request_permission".to_owned(),
5523 params_json: String::new(),
5524 decision: ProviderHostRequestDecision::Denied {
5525 by: ProviderHostDecisionAuthority::Operator
5526 },
5527 outcome: ProviderHostRequestOutcome::Executed,
5528 reason: None,
5529 })
5530 );
5531 }
5532
5533 #[test]
5534 fn unrecognized_notification_reaches_the_operator_as_a_raw_event_instead_of_vanishing() {
5535 let mapped = super::provider_event(
5536 AgentEvent::RpcNotification {
5537 method: "session/update".to_owned(),
5538 params: Default::default(),
5539 },
5540 &[],
5541 );
5542 assert_eq!(
5543 mapped,
5544 Some(ProviderEvent::UnrecognizedNotification {
5545 method: "session/update".to_owned(),
5546 payload_json: "null".to_owned(),
5547 })
5548 );
5549 }
5550
5551 #[test]
5552 fn pty_terminal_capability_defaults_never_override_a_caller_supplied_term() {
5553 let caller_supplied = vec![EnvMutation {
5554 key: OsString::from("TERM"),
5555 value: Some(OsString::from("dumb")),
5556 }];
5557 let filled = with_pty_terminal_capability_defaults(caller_supplied);
5558
5559 let term_values: Vec<_> = filled
5560 .iter()
5561 .filter(|mutation| mutation.key.as_os_str() == OsStr::new("TERM"))
5562 .collect();
5563 assert_eq!(term_values.len(), 1, "TERM must not be duplicated");
5564 assert_eq!(term_values[0].value.as_deref(), Some(OsStr::new("dumb")));
5565
5566 let colorterm = filled
5567 .iter()
5568 .find(|mutation| mutation.key.as_os_str() == OsStr::new("COLORTERM"))
5569 .expect("COLORTERM default is filled in when the caller left it unset");
5570 assert_eq!(colorterm.value.as_deref(), Some(OsStr::new("truecolor")));
5571 }
5572
5573 #[test]
5579 fn mcp_server_acp_entry_translates_a_spec_into_one_stdio_entry() {
5580 let spec = McpServerSpec::new(
5581 "test-mcp-server",
5582 OsString::from("C:\\example\\example-mcp-server.exe"),
5583 vec![OsString::from("--session-proxy")],
5584 vec![
5585 (
5586 OsString::from("TEST_MCP_SESSION_ENDPOINT"),
5587 OsString::from("\\\\.\\pipe\\example-mcp-server-s1"),
5588 ),
5589 (
5590 OsString::from("TEST_MCP_SESSION_TOKEN"),
5591 OsString::from("tok-abc"),
5592 ),
5593 ],
5594 )
5595 .unwrap();
5596
5597 let server = mcp_server_acp_entry(&spec);
5598 match server {
5599 McpServerConfig::Stdio { name, command, args, env } => {
5600 assert_eq!(name, "test-mcp-server", "ACP v1 requires a name; the adapter registers by it");
5601 assert_eq!(command, "C:\\example\\example-mcp-server.exe");
5602 assert_eq!(args, vec!["--session-proxy".to_string()]);
5603 assert_eq!(env.len(), 2, "exactly the entries the spec was given, nothing else");
5604 assert_eq!(
5605 env.iter().find(|pair| pair.name == "TEST_MCP_SESSION_ENDPOINT").map(|pair| pair.value.as_str()),
5606 Some("\\\\.\\pipe\\example-mcp-server-s1")
5607 );
5608 assert_eq!(
5609 env.iter().find(|pair| pair.name == "TEST_MCP_SESSION_TOKEN").map(|pair| pair.value.as_str()),
5610 Some("tok-abc")
5611 );
5612 }
5613 McpServerConfig::Sse { .. } => panic!("expected Stdio variant"),
5614 }
5615 }
5616
5617 #[test]
5618 fn mcp_server_acp_entry_carries_zero_env_pairs_when_the_spec_carries_none() {
5619 let spec = McpServerSpec::new(
5620 "test-mcp-server",
5621 OsString::from("C:\\example\\example-mcp-server.exe"),
5622 vec![OsString::from("--session-proxy")],
5623 Vec::new(),
5624 )
5625 .unwrap();
5626
5627 let server = mcp_server_acp_entry(&spec);
5628 match server {
5629 McpServerConfig::Stdio { env, .. } => assert!(env.is_empty()),
5630 McpServerConfig::Sse { .. } => panic!("expected Stdio variant"),
5631 }
5632 }
5633
5634 #[test]
5635 fn credential_shape_rule_matches_prefixed_hex64_not_bare_digest_or_uuid() {
5636 let hex64 = "a1b2c3d4e5f6".repeat(5) + "a1b2";
5637 assert_eq!(hex64.len(), 64);
5638
5639 let prefixed_token = format!("g4aho_{hex64}");
5640 assert!(
5641 argument_looks_like_credential(&prefixed_token),
5642 "a short lowercase prefix plus '_' plus 64 lowercase hex chars is a credential"
5643 );
5644
5645 assert!(
5646 !argument_looks_like_credential(&hex64),
5647 "a bare 64-hex digest with no prefix is a legitimate diagnostic, not a credential"
5648 );
5649
5650 assert!(
5651 !argument_looks_like_credential("550e8400-e29b-41d4-a716-446655440000"),
5652 "a UUID must not be mistaken for a prefixed-hex64 credential"
5653 );
5654 }
5655
5656 #[test]
5657 fn native_instance_launch_arguments_fail_closed_before_shell_spawn() {
5658 let claude = AgentId::new("claude").unwrap();
5659 assert_eq!(
5660 validate_instance_launch_arguments(
5661 &claude,
5662 TransportKind::Pipe,
5663 &[OsString::from("--bundle-mode")],
5664 )
5665 .unwrap_err(),
5666 "native instance launch arguments require PTY transport"
5667 );
5668 assert_eq!(
5669 validate_instance_launch_arguments(
5670 &claude,
5671 TransportKind::Pty,
5672 &[OsString::from("--resume=session-secret")],
5673 )
5674 .unwrap_err(),
5675 "native instance launch arguments conflict with Claude session, resume, or prompt authority"
5676 );
5677 assert!(validate_instance_launch_arguments(
5678 &claude,
5679 TransportKind::Pty,
5680 &[
5681 OsString::from("--permission-mode"),
5682 OsString::from("default"),
5683 ],
5684 )
5685 .is_ok());
5686 }
5687
5688 fn kind_of(contents: &str) -> Option<OperatorGateKind> {
5692 startup_operator_gate(contents).map(|gate| gate.kind)
5693 }
5694
5695 #[test]
5696 fn startup_operator_gates_are_classified_without_returning_terminal_text() {
5697 assert_eq!(kind_of(" Trust this\nfolder? "), Some(OperatorGateKind::WorkspaceTrust));
5698 assert_eq!(
5699 kind_of("No auth type is selected"),
5700 Some(OperatorGateKind::Authentication)
5701 );
5702 assert_eq!(
5703 kind_of("Sign in with OpenAI to use Codex"),
5704 Some(OperatorGateKind::Authentication)
5705 );
5706 assert_eq!(
5712 kind_of(
5713 "Claude Code can be used with your Claude subscription or billed based on \
5714 API usage through your Console account.\n\
5715 Select login method:\n\
5716 \u{276f} 1. Claude ...\n\
5717 2. API usage billing\n\
5718 3. 3rd-party platform \u{b7} Amazon Bedrock ..."
5719 ),
5720 Some(OperatorGateKind::Authentication)
5721 );
5722 assert_eq!(
5723 kind_of(
5724 "Browser didn't open? Use the url below to sign in (c to copy)\n\
5725 https://claude.com/oauth/authorize?client_id=abc&redirect_uri=https%3A%2F%2Fconsole.anthropic.com&code_challenge=xyz\n\
5726 Paste code here"
5727 ),
5728 Some(OperatorGateKind::Authentication)
5729 );
5730 assert_eq!(
5735 kind_of(
5736 "Once you sign in, copy the generated token and paste it into your .env file; \
5737 no code entry happens on this screen."
5738 ),
5739 None
5740 );
5741 assert_eq!(
5742 kind_of("Kimi Code Update Available\nInstall update now (0.32.0)\nEnter confirm"),
5743 Some(OperatorGateKind::VendorUpdate)
5744 );
5745 assert_eq!(
5746 kind_of(
5747 "Welcome to Claude Code\nChoose the text style that looks best with your terminal"
5748 ),
5749 Some(OperatorGateKind::TerminalAppearance)
5750 );
5751 assert_eq!(
5752 kind_of(
5753 "Welcome to Claude Code for VS Code\nClaude has context of open files and selected lines"
5754 ),
5755 Some(OperatorGateKind::Onboarding)
5756 );
5757 assert_eq!(
5758 kind_of("Welcome to Claude Code\n❯ Press Enter to continue"),
5759 Some(OperatorGateKind::Onboarding)
5760 );
5761 assert_eq!(kind_of("Welcome to Claude Code\n❯ ready\nEnter to send"), None);
5762 assert_eq!(
5763 kind_of(
5764 "Quick safety check: Is this a project you trust?\nYes, I trust this folder\nNo, continue without these permissions"
5765 ),
5766 Some(OperatorGateKind::WorkspaceTrust)
5767 );
5768 assert_eq!(kind_of("ready for a prompt"), None);
5769 assert_eq!(
5774 kind_of(
5775 "Updating Codex via `npm install -g @openai/codex`...\n\
5776 npm warn cleanup Failed to remove some directories\n\
5777 Update ran successfully! Please restart Codex."
5778 ),
5779 Some(OperatorGateKind::VendorUpdate)
5780 );
5781 assert_eq!(
5791 kind_of(
5792 " Hooks need review\n\
5793 6 hooks are new or changed.\n\
5794 Hooks can run outside the sandbox after you trust them.\u{203a} 1. Review hooks\n\
5795 2. Trust all and continue\n\
5796 3. Continue without trusting (hooks won't run) Press enter to confirm or esc to go back"
5797 ),
5798 Some(OperatorGateKind::HookTrust)
5799 );
5800 assert_eq!(
5803 kind_of(
5804 " Hooks need review\n\
5805 42 hooks are new or changed.\n\
5806 Hooks can run outside the sandbox after you trust them.\u{203a} 1. Review hooks\n\
5807 2. Trust all and continue\n\
5808 3. Continue without trusting (hooks won't run) Press enter to confirm or esc to go back"
5809 ),
5810 Some(OperatorGateKind::HookTrust)
5811 );
5812 assert_eq!(
5816 kind_of(
5817 "I reviewed the pre-commit hooks in this repo and they look safe to trust; \
5818 I'll leave the hooks config as-is and continue with the refactor."
5819 ),
5820 None
5821 );
5822 }
5823
5824 #[test]
5831 fn startup_operator_gate_parses_a_numbered_hook_trust_option_list() {
5832 let contents = "Hooks need review\n\
5833 \u{203a} 1. Review hooks\n \
5834 2. Trust all and continue\n \
5835 3. Continue without trusting (hooks won't run)\n\
5836 Press enter to confirm or esc to go back";
5837 let gate = startup_operator_gate(contents).expect("hook trust must classify");
5838 assert_eq!(gate.kind, OperatorGateKind::HookTrust);
5839 assert_eq!(gate.input, OperatorGateInput::NumberedList);
5840 assert_eq!(
5841 gate.options,
5842 vec![
5843 gate4agent_types::OperatorGateOption {
5844 text: "Review hooks".to_owned(),
5845 semantics: OperatorGateOptionSemantics::Inspect,
5846 selected: true,
5847 },
5848 gate4agent_types::OperatorGateOption {
5849 text: "Trust all and continue".to_owned(),
5850 semantics: OperatorGateOptionSemantics::Accept,
5851 selected: false,
5852 },
5853 gate4agent_types::OperatorGateOption {
5854 text: "Continue without trusting (hooks won't run)".to_owned(),
5855 semantics: OperatorGateOptionSemantics::Decline,
5856 selected: false,
5857 },
5858 ],
5859 );
5860 }
5861
5862 #[test]
5869 fn startup_operator_gate_parses_an_arrow_workspace_trust_option_list() {
5870 let contents = "Do you trust the files in this folder?\n \
5871 Trust this folder\n \
5872 Enable project MCP servers. Remembered for this folder.\n\u{276f} \
5873 Don't trust\n \
5874 Exit Kimi Code. Asked again next launch.\n \
5875 \u{2191}\u{2193} navigate \u{b7} Enter select \u{b7} Esc exit";
5876 let gate = startup_operator_gate(contents).expect("workspace trust must classify");
5877 assert_eq!(gate.kind, OperatorGateKind::WorkspaceTrust);
5878 assert_eq!(gate.input, OperatorGateInput::ArrowList);
5879 assert_eq!(
5880 gate.options,
5881 vec![
5882 gate4agent_types::OperatorGateOption {
5883 text: "Trust this folder".to_owned(),
5884 semantics: OperatorGateOptionSemantics::Accept,
5885 selected: false,
5886 },
5887 gate4agent_types::OperatorGateOption {
5888 text: "Don't trust".to_owned(),
5889 semantics: OperatorGateOptionSemantics::Decline,
5890 selected: true,
5891 },
5892 ],
5893 );
5894 }
5895
5896 #[test]
5900 fn startup_operator_gate_without_a_recognized_option_list_reports_unknown_input() {
5901 let gate = startup_operator_gate("No auth type is selected")
5902 .expect("authentication must classify");
5903 assert_eq!(gate.kind, OperatorGateKind::Authentication);
5904 assert_eq!(gate.input, OperatorGateInput::Unknown);
5905 assert!(gate.options.is_empty());
5906 }
5907
5908 #[test]
5917 fn startup_operator_gate_parses_claudes_login_method_chooser_including_its_cursor_row() {
5918 let contents = "Select login method:\n\
5919 \u{276f} 1. Claude ...\n\
5920 2. API usage billing\n\
5921 3. 3rd-party platform \u{b7} Amazon Bedrock ...";
5922 let gate = startup_operator_gate(contents).expect("authentication must classify");
5923 assert_eq!(gate.kind, OperatorGateKind::Authentication);
5924 assert_eq!(gate.subject, OperatorGateSubject::Account);
5925 assert_eq!(gate.input, OperatorGateInput::NumberedList);
5926 assert_eq!(
5927 gate.options,
5928 vec![
5929 gate4agent_types::OperatorGateOption {
5930 text: "Claude ...".to_owned(),
5931 semantics: OperatorGateOptionSemantics::Unknown,
5932 selected: true,
5933 },
5934 gate4agent_types::OperatorGateOption {
5935 text: "API usage billing".to_owned(),
5936 semantics: OperatorGateOptionSemantics::Unknown,
5937 selected: false,
5938 },
5939 gate4agent_types::OperatorGateOption {
5940 text: "3rd-party platform \u{b7} Amazon Bedrock ...".to_owned(),
5941 semantics: OperatorGateOptionSemantics::Unknown,
5942 selected: false,
5943 },
5944 ],
5945 );
5946 }
5947
5948 #[test]
5953 fn startup_operator_gate_oauth_wait_screen_reports_text_entry_input() {
5954 let contents = "Browser didn't open? Use the url below to sign in (c to copy)\n\
5955 https://claude.com/oauth/authorize?client_id=abc&redirect_uri=https%3A%2F%2Fconsole.anthropic.com&code_challenge=xyz\n\
5956 Paste code here";
5957 let gate = startup_operator_gate(contents).expect("oauth wait screen must classify");
5958 assert_eq!(gate.kind, OperatorGateKind::Authentication);
5959 assert_eq!(gate.subject, OperatorGateSubject::Account);
5960 assert_eq!(gate.input, OperatorGateInput::TextEntry);
5961 assert!(gate.options.is_empty());
5962 }
5963
5964 #[test]
5970 fn parse_operator_gate_options_recognizes_numbered_and_arrow_shapes_and_nothing_else() {
5971 assert_eq!(
5972 parse_operator_gate_options("Press enter to confirm or esc to go back"),
5973 (OperatorGateInput::Unknown, Vec::new()),
5974 );
5975 let (input, options) = parse_operator_gate_options(
5976 "\u{203a} 1. Review hooks\n 2. Trust all and continue",
5977 );
5978 assert_eq!(input, OperatorGateInput::NumberedList);
5979 assert_eq!(options.len(), 2);
5980 let (input, options) = parse_operator_gate_options(
5981 " Trust this folder\n Enable project MCP servers. Remembered for this folder.\n\u{276f} Don't trust",
5982 );
5983 assert_eq!(input, OperatorGateInput::ArrowList);
5984 assert_eq!(options.len(), 2);
5985 }
5986
5987 #[test]
5988 fn classify_operator_gate_option_semantics_reads_the_verb_not_the_position() {
5989 assert_eq!(
5990 classify_operator_gate_option_semantics("Review hooks"),
5991 OperatorGateOptionSemantics::Inspect,
5992 );
5993 assert_eq!(
5994 classify_operator_gate_option_semantics("Trust all and continue"),
5995 OperatorGateOptionSemantics::Accept,
5996 );
5997 assert_eq!(
5998 classify_operator_gate_option_semantics("Continue without trusting (hooks won't run)"),
5999 OperatorGateOptionSemantics::Decline,
6000 );
6001 assert_eq!(
6002 classify_operator_gate_option_semantics("Don't trust"),
6003 OperatorGateOptionSemantics::Decline,
6004 );
6005 assert_eq!(
6006 classify_operator_gate_option_semantics("Exit Kimi Code"),
6007 OperatorGateOptionSemantics::Exit,
6008 );
6009 assert_eq!(
6010 classify_operator_gate_option_semantics("Something unrecognized"),
6011 OperatorGateOptionSemantics::Unknown,
6012 );
6013 }
6014
6015 #[test]
6016 fn screen_failure_recognizes_crash_and_missing_command_banners() {
6017 assert_eq!(
6018 screen_failure("thread 'main' panicked at 'index out of bounds', src/main.rs:12:5"),
6019 Some("crash")
6020 );
6021 assert_eq!(
6022 screen_failure("Traceback (most recent call last):\n File \"a.py\", line 1"),
6023 Some("crash")
6024 );
6025 assert_eq!(
6026 screen_failure("bash: fooagent: command not found"),
6027 Some("missing command")
6028 );
6029 assert_eq!(
6030 screen_failure("C:\\workspace> fooagent\n'fooagent' is not recognized as an internal or external command"),
6031 Some("missing command")
6032 );
6033 assert_eq!(
6034 screen_failure(
6035 "workspace/project $ fooagent: no such file or directory\nworkspace/project $"
6036 ),
6037 Some("missing command")
6038 );
6039 assert_eq!(
6042 screen_failure("The build log mentions no such file or directory near line 40."),
6043 None
6044 );
6045 assert_eq!(
6049 screen_failure("fooagent-installer: fatal error: missing header <stdio.h>"),
6050 None
6051 );
6052 }
6053
6054 #[test]
6055 fn screen_failure_recognizes_grok_missing_api_key_banner() {
6056 assert_eq!(
6061 screen_failure(
6062 "\u{274c} Error: API key required. Set GROK_API_KEY environment variable, \
6063 use --api-key flag, or set \"apiKey\" field in ~/.grok/user-settings.json"
6064 ),
6065 Some("missing api key")
6066 );
6067 assert_eq!(
6071 screen_failure(
6072 "You'll need an API key for this provider before it works; check the docs \
6073 for how to configure one."
6074 ),
6075 None
6076 );
6077 }
6078
6079 #[test]
6080 fn screen_failure_no_longer_has_an_authentication_expired_bucket() {
6081 assert_eq!(
6086 screen_failure("Your session token has expired. Please re-authenticate."),
6087 None
6088 );
6089 assert_eq!(
6090 screen_failure("The stored credential was revoked by the workspace admin."),
6091 None
6092 );
6093 }
6094
6095 #[test]
6096 fn screen_failure_does_not_false_positive_on_an_agent_narrating_an_error() {
6097 let narration = "I looked at the traceback you pasted -- it's a Python \
6102 exception, a plain ValueError from a retry loop, not a real crash. \
6103 Our harness logs the word fatal in its own banner for visibility, \
6104 but nothing actually panicked. The command definitely exists; it \
6105 just needed different flags, and your token is still valid.";
6106 assert_eq!(screen_failure(narration), None);
6107 }
6108
6109 #[test]
6110 fn screen_failure_for_generation_ignores_every_marker_once_the_session_was_ever_ready() {
6111 let shell_command_not_found = "$ frobnicate --help\nbash: frobnicate: command not found\n$";
6116 assert_eq!(
6117 screen_failure_for_generation(shell_command_not_found, true),
6118 None
6119 );
6120
6121 let python_traceback_from_a_ran_script = "$ python broken.py\n\
6122 Traceback (most recent call last):\n File \"broken.py\", line 3, in <module>\n\
6123 ValueError: bad input\n$";
6124 assert_eq!(
6125 screen_failure_for_generation(python_traceback_from_a_ran_script, true),
6126 None
6127 );
6128
6129 let cargo_test_panic = "running 1 test\n\
6130 error[E0308]: mismatched types\n\
6131 thread 'tests::it_fails' panicked at 'assertion failed', src/lib.rs:9:5\n\
6132 test result: FAILED. 0 passed; 1 failed";
6133 assert_eq!(screen_failure_for_generation(cargo_test_panic, true), None);
6134
6135 let agent_narrating_an_expired_token = "The API token has expired; \
6136 I revoked the old credential and issued a new one for you.";
6137 assert_eq!(
6138 screen_failure_for_generation(agent_narrating_an_expired_token, true),
6139 None
6140 );
6141
6142 assert_eq!(
6147 screen_failure_for_generation(shell_command_not_found, false),
6148 Some("missing command")
6149 );
6150 assert_eq!(
6151 screen_failure_for_generation(python_traceback_from_a_ran_script, false),
6152 Some("crash")
6153 );
6154 assert_eq!(screen_failure_for_generation(cargo_test_panic, false), Some("crash"));
6155 }
6156
6157 #[test]
6158 fn classify_pty_screen_state_lets_a_foreign_process_win_over_a_matching_gate_text() {
6159 let foreign = ForegroundVerdict::Foreign {
6164 process: "npm".to_owned(),
6165 };
6166 assert_eq!(
6167 classify_pty_screen_state(Some(&foreign), None, None),
6168 PtyScreenState::NotAgent {
6169 observed_process: "npm".to_owned()
6170 }
6171 );
6172 let vendor_update_gate = OperatorGateState::new(OperatorGateKind::VendorUpdate);
6173 assert_eq!(
6174 classify_pty_screen_state(Some(&foreign), Some(&vendor_update_gate), None),
6175 PtyScreenState::NotAgent {
6176 observed_process: "npm".to_owned()
6177 }
6178 );
6179 }
6180
6181 #[test]
6182 fn classify_pty_screen_state_reports_a_gate_only_when_foreground_matches() {
6183 let authentication_gate = OperatorGateState::new(OperatorGateKind::Authentication);
6184 assert_eq!(
6185 classify_pty_screen_state(Some(&ForegroundVerdict::Agent), Some(&authentication_gate), None),
6186 PtyScreenState::OperatorGate {
6187 gate: authentication_gate.clone()
6188 }
6189 );
6190 }
6191
6192 #[test]
6193 fn classify_pty_screen_state_reports_failing_for_a_matched_foreground() {
6194 assert_eq!(
6195 classify_pty_screen_state(Some(&ForegroundVerdict::Agent), None, Some("crash")),
6196 PtyScreenState::Failing {
6197 reason: "crash".to_owned()
6198 }
6199 );
6200 }
6201
6202 #[test]
6203 fn classify_pty_screen_state_prefers_failing_over_a_co_occurring_gate() {
6204 let authentication_gate = OperatorGateState::new(OperatorGateKind::Authentication);
6208 assert_eq!(
6209 classify_pty_screen_state(
6210 Some(&ForegroundVerdict::Agent),
6211 Some(&authentication_gate),
6212 Some("crash")
6213 ),
6214 PtyScreenState::Failing {
6215 reason: "crash".to_owned()
6216 }
6217 );
6218 }
6219
6220 #[test]
6221 fn classify_pty_screen_state_reaches_ready_only_with_matched_foreground_and_clean_text() {
6222 assert_eq!(
6223 classify_pty_screen_state(Some(&ForegroundVerdict::Agent), None, None),
6224 PtyScreenState::Ready
6225 );
6226 }
6227
6228 #[test]
6229 fn classify_pty_screen_state_never_reaches_ready_without_a_foreground_observation() {
6230 assert_eq!(classify_pty_screen_state(None, None, None), PtyScreenState::Unknown);
6233 }
6234
6235 #[test]
6236 fn resolve_foreground_verdict_tolerates_a_spawning_wrapper_like_readiness_does() {
6237 let spec = gate4agent_testkit::interactive_agent_spec();
6245 let observation = PtyForegroundObservation {
6246 root_pid: 1,
6247 observed_pid: 2,
6248 observed_process: "node".to_owned(),
6249 readiness: ForegroundObservation {
6250 process_name: Some("node".to_owned()),
6251 has_child_processes: true,
6252 is_shell: false,
6253 },
6254 source: PtyForegroundSource::ProcessTree,
6255 };
6256 let verdict = resolve_foreground_verdict(&spec, &observation, RuntimePlatform::current());
6257 assert_eq!(verdict, ForegroundVerdict::Agent);
6258 assert_eq!(
6259 classify_pty_screen_state(Some(&verdict), None, None),
6260 PtyScreenState::Ready
6261 );
6262 }
6263
6264 #[test]
6265 fn foreground_probe_disarms_only_once_ready_and_rearms_on_a_fresh_gate() {
6266 assert_eq!(
6267 foreground_probe_schedule(&PtyScreenState::Ready),
6268 ForegroundProbeSchedule::Disarmed
6269 );
6270 assert_eq!(
6271 foreground_probe_schedule(&PtyScreenState::Unknown),
6272 ForegroundProbeSchedule::Armed
6273 );
6274 assert_eq!(
6275 foreground_probe_schedule(&PtyScreenState::OperatorGate {
6276 gate: OperatorGateState::new(OperatorGateKind::VendorUpdate)
6277 }),
6278 ForegroundProbeSchedule::Armed
6279 );
6280 }
6281
6282 #[test]
6283 fn text_only_gate_transition_rearms_the_probe_only_when_leaving_ready() {
6284 let gate = PtyScreenState::OperatorGate {
6285 gate: OperatorGateState::new(OperatorGateKind::VendorUpdate),
6286 };
6287 assert!(foreground_probe_rearms_immediately(
6288 &PtyScreenState::Ready,
6289 &gate
6290 ));
6291 assert!(!foreground_probe_rearms_immediately(
6294 &PtyScreenState::Unknown,
6295 &gate
6296 ));
6297 }
6298
6299 #[test]
6300 fn readiness_diagnostics_detects_an_ansi_split_gate_without_exposing_text() {
6301 let mut diagnostics = ReadinessDiagnostics::default();
6302 diagnostics.observe_output(b"\x1b[31mNo auth ");
6303 diagnostics.observe_output(b"\x1b[0mtype is selected");
6304 assert_eq!(
6305 diagnostics.operator_gate.as_ref().map(|gate| gate.kind),
6306 Some(OperatorGateKind::Authentication)
6307 );
6308 assert!(!diagnostics.summary().contains("No auth"));
6309 }
6310
6311 #[test]
6312 fn semantic_utf8_decoder_preserves_codepoints_split_across_pty_reads() {
6313 let mut decoder = Utf8ChunkDecoder::default();
6314 let bytes = "ready Привет".as_bytes();
6315 let split = bytes
6316 .windows(2)
6317 .position(|window| window[0] >= 0x80 && window[1] >= 0x80)
6318 .expect("Cyrillic text contains adjacent UTF-8 bytes")
6319 + 1;
6320 let first = decoder.push(&bytes[..split]);
6321 let second = decoder.push(&bytes[split..]);
6322 assert_eq!(format!("{first}{second}"), "ready Привет");
6323 assert!(!first.contains('\u{fffd}'));
6324 assert!(!second.contains('\u{fffd}'));
6325 }
6326
6327 #[test]
6328 fn semantic_utf8_decoder_replaces_invalid_bytes_without_stalling() {
6329 let mut decoder = Utf8ChunkDecoder::default();
6330 assert_eq!(decoder.push(b"ok\xfftail"), "ok\u{fffd}tail");
6331 assert_eq!(decoder.push(" Привет".as_bytes()), " Привет");
6332 }
6333
6334 #[test]
6335 fn fresh_claude_pty_preassigns_the_exact_vendor_session_id_argv() {
6336 let claude = builtin_adapter_registry()
6337 .binding(AdapterFamily::PtySemantic, "claude-code")
6338 .expect("Claude PTY binding");
6339 let mut args = Vec::new();
6340 let identity = prepare_fresh_pty_provider_session(Some(claude), false, true, &mut args)
6341 .expect("fresh Claude provider identity");
6342 let parsed = uuid::Uuid::parse_str(&identity.id).expect("valid Claude UUID");
6343 assert_eq!(parsed.get_version_num(), 4);
6344 assert_eq!(identity.key, gate4agent_types::ProviderSessionKey::SessionId);
6345 assert!(identity.transcript_path.is_none());
6346 assert_eq!(args, [OsString::from("--session-id"), OsString::from(identity.id)]);
6347 }
6348
6349 #[test]
6350 fn fresh_codex_and_resumed_claude_do_not_preassign_a_second_identity() {
6351 let adapters = builtin_adapter_registry();
6352 let codex = adapters
6353 .binding(AdapterFamily::PtySemantic, "codex")
6354 .expect("Codex PTY binding");
6355 let claude = adapters
6356 .binding(AdapterFamily::PtySemantic, "claude-code")
6357 .expect("Claude PTY binding");
6358 let mut codex_args = Vec::new();
6359 assert!(prepare_fresh_pty_provider_session(Some(codex), false, true, &mut codex_args)
6360 .is_none());
6361 assert!(codex_args.is_empty());
6362
6363 let mut resume_args = vec![OsString::from("--resume"), OsString::from("vendor-id")];
6364 assert!(prepare_fresh_pty_provider_session(Some(claude), true, true, &mut resume_args)
6365 .is_none());
6366 assert_eq!(
6367 resume_args,
6368 [OsString::from("--resume"), OsString::from("vendor-id")]
6369 );
6370 }
6371
6372 #[test]
6373 fn raw_pty_policy_omits_all_identity_and_semantic_startup_paths() {
6374 let adapters = builtin_adapter_registry();
6375 let claude = adapters
6376 .binding(AdapterFamily::PtySemantic, "claude-code")
6377 .expect("Claude PTY binding");
6378 let codex = adapters
6379 .binding(AdapterFamily::PtySemantic, "codex")
6380 .expect("Codex PTY binding");
6381 let kimi = adapters
6382 .binding(AdapterFamily::PtySemantic, "kimi")
6383 .expect("Kimi PTY binding");
6384 let policy = ProviderRuntimePolicy::raw_pty();
6385 let mut args = Vec::new();
6386
6387 assert!(prepare_fresh_pty_provider_session(Some(claude), false, false, &mut args)
6388 .is_none());
6389 assert!(args.is_empty());
6390 assert!(!should_probe_pty_identity(policy, Some(codex), false, "codex"));
6391 assert!(!should_probe_pty_identity(policy, Some(kimi), false, "kimi"));
6392 assert!(!should_attach_pty_provider_stream(policy));
6393 }
6394
6395 #[test]
6396 fn identity_probe_requires_structured_prompt_for_codex_and_kimi() {
6397 let adapters = builtin_adapter_registry();
6398 let codex = adapters
6399 .binding(AdapterFamily::PtySemantic, "codex")
6400 .expect("Codex PTY binding");
6401 let kimi = adapters
6402 .binding(AdapterFamily::PtySemantic, "kimi")
6403 .expect("Kimi PTY binding");
6404 let without_structured_prompt =
6405 ProviderRuntimePolicy::new(true, true, false, true, false, false)
6406 .expect("identity observation policy without structured prompt");
6407
6408 assert!(!should_probe_pty_identity(
6409 without_structured_prompt,
6410 Some(codex),
6411 false,
6412 "codex",
6413 ));
6414 assert!(!should_probe_pty_identity(
6415 without_structured_prompt,
6416 Some(kimi),
6417 false,
6418 "kimi",
6419 ));
6420
6421 let with_structured_prompt =
6422 ProviderRuntimePolicy::new(true, true, true, true, false, false)
6423 .expect("identity probe policy");
6424 assert!(should_probe_pty_identity(
6425 with_structured_prompt,
6426 Some(codex),
6427 false,
6428 "codex",
6429 ));
6430 assert!(should_probe_pty_identity(
6431 with_structured_prompt,
6432 Some(kimi),
6433 false,
6434 "kimi",
6435 ));
6436 }
6437
6438 #[test]
6439 fn runtime_policy_admits_raw_native_resume_but_denies_unverified_prompt_injection() {
6440 let raw = ProviderRuntimePolicy::raw_pty();
6441 assert!(validate_spawn_runtime_policy(raw, TransportKind::Pty, false, false).is_ok());
6442 assert!(validate_spawn_runtime_policy(raw, TransportKind::Pty, true, false)
6443 .unwrap_err()
6444 .contains("SemanticReadiness"));
6445 assert!(validate_spawn_runtime_policy(raw, TransportKind::Pty, false, true).is_ok());
6446 assert!(validate_spawn_runtime_policy(raw, TransportKind::Pty, true, true)
6447 .unwrap_err()
6448 .contains("SemanticReadiness"));
6449
6450 let resume_without_prompt =
6451 ProviderRuntimePolicy::new(true, false, false, true, true, false)
6452 .expect("identity and resume policy");
6453 assert!(validate_spawn_runtime_policy(
6454 resume_without_prompt,
6455 TransportKind::Pty,
6456 false,
6457 true,
6458 )
6459 .is_ok());
6460 assert!(validate_spawn_runtime_policy(
6461 resume_without_prompt,
6462 TransportKind::Pty,
6463 true,
6464 true,
6465 )
6466 .unwrap_err()
6467 .contains("SemanticReadiness"));
6468 }
6469
6470 #[test]
6471 fn acp_transport_bypasses_the_pty_semantic_policy_gate_pty_and_pipe_still_enforce_it() {
6472 let no_pty_capabilities_at_all = ProviderRuntimePolicy::new(
6478 false, false, false, false, false, false,
6479 )
6480 .expect("an all-false policy is internally valid");
6481 assert!(validate_spawn_runtime_policy(
6482 no_pty_capabilities_at_all,
6483 TransportKind::Acp,
6484 false,
6485 false,
6486 )
6487 .is_ok());
6488 assert!(validate_spawn_runtime_policy(
6492 no_pty_capabilities_at_all,
6493 TransportKind::Pty,
6494 false,
6495 false,
6496 )
6497 .unwrap_err()
6498 .contains("RawPtyLifecycle"));
6499 assert!(validate_spawn_runtime_policy(
6500 no_pty_capabilities_at_all,
6501 TransportKind::Pipe,
6502 false,
6503 false,
6504 )
6505 .unwrap_err()
6506 .contains("RawPtyLifecycle"));
6507 }
6508
6509 #[test]
6510 fn prompt_render_probe_ignores_terminal_wrapping_and_uses_the_tail() {
6511 let prompt = "prefix with spaces\nand punctuation: final-render-token-1234567890";
6512 let probe = prompt_render_probe(prompt);
6513 assert!("screen prefix with spaces and punctuation final render token 1234567890"
6514 .chars()
6515 .filter(|character| character.is_alphanumeric())
6516 .flat_map(char::to_lowercase)
6517 .collect::<String>()
6518 .contains(&probe));
6519 assert!(probe.chars().count() <= 32);
6520 }
6521
6522 #[test]
6523 fn prompt_render_requires_new_sequence_and_visible_tail_evidence() {
6524 let probe = prompt_render_probe("final render token");
6525 let baseline = snapshot(7, "old composer");
6526 assert!(!prompt_rendered(
6527 &snapshot(7, "final\nrender\ntoken"),
6528 &baseline,
6529 &probe
6530 ));
6531 assert!(!prompt_rendered(
6532 &snapshot(8, "unrelated redraw"),
6533 &baseline,
6534 &probe
6535 ));
6536 assert!(prompt_rendered(
6537 &snapshot(8, "composer\nfinal\nrender\ntoken"),
6538 &baseline,
6539 &probe
6540 ));
6541 assert!(prompt_rendered(
6542 &snapshot(8, "composer [Pasted Content 4096 chars]"),
6543 &baseline,
6544 &probe
6545 ));
6546 assert!(prompt_rendered(
6547 &snapshot(8, "composer [Pasted text #1 +6 lines]"),
6548 &baseline,
6549 &probe
6550 ));
6551 }
6552
6553 #[test]
6554 fn an_existing_paste_placeholder_cannot_pass_on_an_unrelated_redraw() {
6555 let probe = prompt_render_probe("a new long prompt");
6556 assert!(!prompt_rendered(
6557 &snapshot(8, "composer [Pasted Content 4096 chars]\nunrelated redraw"),
6558 &snapshot(7, "composer [Pasted Content 4096 chars]"),
6559 &probe
6560 ));
6561 assert!(!prompt_rendered(
6562 &snapshot(8, "composer [Pasted text #1 +6 lines]\nunrelated redraw"),
6563 &snapshot(7, "composer [Pasted text #1 +6 lines]"),
6564 &probe
6565 ));
6566 }
6567
6568 #[test]
6569 fn punctuation_only_prompt_cannot_pass_on_an_unrelated_redraw() {
6570 let probe = prompt_render_probe("!?---");
6571 assert!(probe.is_empty());
6572 assert!(!prompt_rendered(
6573 &snapshot(2, "unrelated redraw"),
6574 &snapshot(1, "old composer"),
6575 &probe
6576 ));
6577 }
6578
6579 #[test]
6580 fn prompt_probe_matches_the_sanitized_terminal_payload() {
6581 let prompt = "payload\u{1b}tail";
6582 let sanitized = gate4agent_types::sanitize_prompt_text(prompt);
6583 let probe = prompt_render_probe(&sanitized);
6584 assert!(prompt_rendered(
6585 &snapshot(2, "composer payload<ESC>tail"),
6586 &snapshot(1, "old composer"),
6587 &probe
6588 ));
6589 }
6590
6591 #[test]
6592 fn lag_and_data_gap_reserve_missed_provider_source_positions() {
6593 let mut next = 7;
6594 assert_eq!(reserve_provider_gap_sequence(&mut next, 3), Some(9));
6595 assert_eq!(next, 10);
6596 assert_eq!(reserve_provider_gap_sequence(&mut next, 0), None);
6597 assert_eq!(next, 10);
6598
6599 next = u64::MAX;
6600 assert_eq!(reserve_provider_gap_sequence(&mut next, 1), None);
6601 assert_eq!(next, u64::MAX);
6602 }
6603
6604 const CODEX_STATUS_RAW_ANSI: &str = "│ 5h limit: \x1b[m[████████████████████] 100% left\x1b[2m (resets 00:08 on 31 Aug) │\r\n│ Weekly limit: \x1b[m[████████████████████] 100% left\x1b[2m (resets 19:08 on 6 Sep) │";
6617
6618 #[test]
6619 fn bare_detector_misses_the_live_ansi_sample_that_rate_limit_feed_catches() {
6620 let bare = RateLimitDetector::new_for_tool(CliTool::Codex);
6623 assert!(bare.detect(CODEX_STATUS_RAW_ANSI).is_none());
6624
6625 let mut feed = RateLimitFeed::new_for_tool(CliTool::Codex);
6628 let info = feed
6629 .detect(CODEX_STATUS_RAW_ANSI)
6630 .expect("quota-state line must be recognized once ANSI is stripped");
6631 assert_eq!(info.limit_type, RateLimitType::Session);
6632 assert_eq!(info.usage_percent, Some(0.0));
6633 assert_eq!(info.resets_at_text.as_deref(), Some("00:08 on 31 Aug"));
6634 }
6635
6636 #[test]
6637 fn rate_limit_feed_reassembles_a_quota_state_line_split_mid_escape_and_mid_row() {
6638 let first =
6642 "│ 5h limit: \x1b[m[████████████████████] 100% left\x1b";
6643 let second = "[2m (resets 00:08 on 31 Aug) │\r\n";
6644
6645 let mut feed = RateLimitFeed::new_for_tool(CliTool::Codex);
6646 assert!(
6647 feed.detect(first).is_none(),
6648 "the row is not complete yet, so nothing should match on the first chunk"
6649 );
6650 let info = feed
6651 .detect(second)
6652 .expect("the row completes once the second chunk arrives");
6653 assert_eq!(info.limit_type, RateLimitType::Session);
6654 assert_eq!(info.resets_at_text.as_deref(), Some("00:08 on 31 Aug"));
6655 }
6656
6657 #[test]
6658 fn rate_limit_feed_strips_ansi_before_matching_a_colored_failure_message() {
6659 let raw = "\x1b[31mError: rate limit exceeded\x1b[0m. Please wait.";
6663 let mut feed = RateLimitFeed::new_for_tool(CliTool::Codex);
6664 let info = feed
6665 .detect(raw)
6666 .expect("a colored refusal message must still be recognized");
6667 assert_eq!(info.limit_type, RateLimitType::Unknown);
6668 }
6669
6670 #[test]
6671 fn rate_limit_feed_bounds_its_buffer_across_a_newline_free_redraw() {
6672 let mut feed = RateLimitFeed::new_for_tool(CliTool::Codex);
6677 let noise = "x".repeat(super::RATE_LIMIT_FEED_BUFFER_MAX_BYTES * 3);
6678 assert!(feed.detect(&noise).is_none());
6679 assert!(feed.buffer.len() <= super::RATE_LIMIT_FEED_BUFFER_MAX_BYTES);
6680
6681 let info = feed
6682 .detect(CODEX_STATUS_RAW_ANSI)
6683 .expect("a real quota-state line must still be recognized after the noise");
6684 assert_eq!(info.limit_type, RateLimitType::Session);
6685 }
6686
6687 #[test]
6692 fn full_auto_maps_to_the_permissive_host_policy() {
6693 assert_eq!(host_policy_for_approval_level(ApprovalLevel::FullAuto), HostPolicy::Yolo);
6697 }
6698
6699 #[test]
6700 fn read_only_maps_to_the_read_only_host_policy() {
6701 assert_eq!(
6702 host_policy_for_approval_level(ApprovalLevel::ReadOnly),
6703 HostPolicy::ReadOnly
6704 );
6705 }
6706
6707 #[test]
6708 fn moderate_and_unmanaged_both_fall_back_to_the_default_host_policy() {
6709 assert_eq!(host_policy_for_approval_level(ApprovalLevel::Moderate), HostPolicy::Auto);
6710 assert_eq!(host_policy_for_approval_level(ApprovalLevel::Unmanaged), HostPolicy::Auto);
6711 assert_eq!(HostPolicy::default(), HostPolicy::Auto);
6712 }
6713
6714 #[test]
6715 fn no_approval_level_ever_maps_to_a_refusing_policy_except_read_only_itself() {
6716 for level in [
6721 ApprovalLevel::FullAuto,
6722 ApprovalLevel::Moderate,
6723 ApprovalLevel::ReadOnly,
6724 ApprovalLevel::Unmanaged,
6725 ] {
6726 assert_ne!(
6727 host_policy_for_approval_level(level),
6728 HostPolicy::Deny,
6729 "{level:?} must never map to HostPolicy::Deny"
6730 );
6731 }
6732 }
6733
6734 #[test]
6742 fn defers_permission_requests_matches_the_verified_provider_table() {
6743 let claude = AgentId::new("claude").unwrap();
6748 assert!(!defers_permission_requests(&claude, ApprovalLevel::FullAuto));
6749 assert!(defers_permission_requests(&claude, ApprovalLevel::Moderate));
6750 assert!(defers_permission_requests(&claude, ApprovalLevel::ReadOnly));
6754 assert!(!defers_permission_requests(&claude, ApprovalLevel::Unmanaged));
6758
6759 let codex = AgentId::new("codex").unwrap();
6760 assert!(!defers_permission_requests(&codex, ApprovalLevel::FullAuto));
6761 assert!(!defers_permission_requests(&codex, ApprovalLevel::Moderate));
6766 assert!(defers_permission_requests(&codex, ApprovalLevel::ReadOnly));
6767 assert!(!defers_permission_requests(&codex, ApprovalLevel::Unmanaged));
6769
6770 let grok = AgentId::new("grok").unwrap();
6771 assert!(!defers_permission_requests(&grok, ApprovalLevel::FullAuto));
6772 assert!(defers_permission_requests(&grok, ApprovalLevel::Moderate));
6773 assert!(defers_permission_requests(&grok, ApprovalLevel::ReadOnly));
6777 assert!(!defers_permission_requests(&grok, ApprovalLevel::Unmanaged));
6779
6780 let kimi = AgentId::new("kimi").unwrap();
6781 assert!(!defers_permission_requests(&kimi, ApprovalLevel::FullAuto));
6782 assert!(!defers_permission_requests(&kimi, ApprovalLevel::Moderate));
6785 assert!(!defers_permission_requests(&kimi, ApprovalLevel::ReadOnly));
6789 assert!(!defers_permission_requests(&kimi, ApprovalLevel::Unmanaged));
6791 }
6792
6793 fn session_mode(id: &str) -> SessionMode {
6800 SessionMode { id: id.to_owned(), name: id.to_owned(), description: None }
6801 }
6802
6803 #[test]
6804 fn required_acp_mode_matches_the_sourced_table() {
6805 let claude = AgentId::new("claude").unwrap();
6810 assert_eq!(
6811 required_acp_mode(&claude, ApprovalLevel::FullAuto),
6812 Ok(Some(ModeId::new("bypassPermissions")))
6813 );
6814 assert_eq!(
6815 required_acp_mode(&claude, ApprovalLevel::Moderate),
6816 Ok(Some(ModeId::new("acceptEdits")))
6817 );
6818 assert_eq!(
6822 required_acp_mode(&claude, ApprovalLevel::ReadOnly),
6823 Ok(Some(ModeId::new("default")))
6824 );
6825 assert_eq!(required_acp_mode(&claude, ApprovalLevel::Unmanaged), Ok(None));
6826
6827 let codex = AgentId::new("codex").unwrap();
6828 assert_eq!(
6829 required_acp_mode(&codex, ApprovalLevel::FullAuto),
6830 Ok(Some(ModeId::new("agent-full-access")))
6831 );
6832 assert_eq!(
6833 required_acp_mode(&codex, ApprovalLevel::Moderate),
6834 Ok(Some(ModeId::new("agent")))
6835 );
6836 assert_eq!(
6837 required_acp_mode(&codex, ApprovalLevel::ReadOnly),
6838 Ok(Some(ModeId::new("read-only")))
6839 );
6840 assert_eq!(required_acp_mode(&codex, ApprovalLevel::Unmanaged), Ok(None));
6841
6842 let kimi = AgentId::new("kimi").unwrap();
6854 assert_eq!(required_acp_mode(&kimi, ApprovalLevel::FullAuto), Ok(None));
6855 assert_eq!(
6856 acp_approval_level_args(&kimi, ApprovalLevel::FullAuto),
6857 vec!["--auto".to_owned()],
6858 "kimi's FullAuto must reach the process through argv, since it has no ACP mode"
6859 );
6860 for level in [ApprovalLevel::Moderate, ApprovalLevel::ReadOnly] {
6861 assert!(
6862 required_acp_mode(&kimi, level).is_err(),
6863 "kimi at {level:?} has neither a sourced ACP mode id nor vendor flags, so it must refuse"
6864 );
6865 }
6866 assert_eq!(required_acp_mode(&kimi, ApprovalLevel::Unmanaged), Ok(None));
6867 assert!(
6868 acp_approval_level_args(&kimi, ApprovalLevel::Unmanaged).is_empty(),
6869 "Unmanaged imposes nothing, so it must add no argv flags either"
6870 );
6871
6872 let grok = AgentId::new("grok").unwrap();
6879 assert_eq!(required_acp_mode(&grok, ApprovalLevel::FullAuto), Ok(None));
6880 assert_eq!(
6881 acp_approval_level_args(&grok, ApprovalLevel::FullAuto),
6882 vec!["--always-approve".to_owned()],
6883 );
6884 assert_eq!(required_acp_mode(&grok, ApprovalLevel::Moderate), Ok(None));
6885 assert_eq!(
6886 acp_approval_level_args(&grok, ApprovalLevel::Moderate),
6887 vec!["--permission-mode".to_owned(), "auto".to_owned()],
6888 );
6889 assert!(
6890 required_acp_mode(&grok, ApprovalLevel::ReadOnly).is_err(),
6891 "grok at ReadOnly is Unsupported outright, so it must refuse"
6892 );
6893 assert_eq!(required_acp_mode(&grok, ApprovalLevel::Unmanaged), Ok(None));
6894 }
6895
6896 #[test]
6907 fn required_acp_mode_never_silently_permits_auto_except_for_unmanaged() {
6908 for id in ["claude", "codex", "grok", "kimi"] {
6909 let agent = AgentId::new(id).unwrap();
6910 for level in [ApprovalLevel::FullAuto, ApprovalLevel::Moderate, ApprovalLevel::ReadOnly] {
6911 match required_acp_mode(&agent, level) {
6912 Ok(Some(mode_id)) => assert_ne!(
6913 mode_id.as_str(),
6914 "auto",
6915 "{id} at {level:?} must never silently resolve the vendor's own default 'auto'"
6916 ),
6917 Ok(None) => assert!(
6918 !acp_approval_level_args(&agent, level).is_empty(),
6919 "{id} at {level:?} resolved Ok(None) while carrying no vendor flags -- \
6920 that leaves the session at the agent's own default. Only Unmanaged, or a \
6921 level argv actually applies, may resolve to no mode"
6922 ),
6923 Err(_) => {} }
6925 }
6926 assert_eq!(
6927 required_acp_mode(&agent, ApprovalLevel::Unmanaged),
6928 Ok(None),
6929 "{id}: Unmanaged is the one level allowed to apply nothing"
6930 );
6931 }
6932 }
6933
6934 #[test]
6935 fn approval_level_not_offered_message_refuses_by_name_and_lists_what_was_offered() {
6936 let offered = vec![session_mode("bypassPermissions"), session_mode("default")];
6937 let mode_id = ModeId::new("plan");
6938 let message = approval_level_not_offered_message(ApprovalLevel::ReadOnly, &mode_id, &offered)
6939 .expect("'plan' was not among the offered modes");
6940 assert!(message.contains("ApprovalLevelNotOfferedByAgent"));
6941 assert!(message.contains("bypassPermissions"));
6942 assert!(message.contains("default"));
6943 assert!(message.contains("ReadOnly"));
6944 }
6945
6946 #[test]
6947 fn approval_level_not_offered_message_is_none_when_the_agent_announced_it() {
6948 let offered = vec![session_mode("plan")];
6949 let mode_id = ModeId::new("plan");
6950 assert_eq!(
6951 approval_level_not_offered_message(ApprovalLevel::ReadOnly, &mode_id, &offered),
6952 None
6953 );
6954 }
6955
6956 #[test]
6957 fn approval_level_not_offered_message_handles_an_agent_that_announces_no_modes_at_all() {
6958 let mode_id = ModeId::new("bypassPermissions");
6959 let message = approval_level_not_offered_message(ApprovalLevel::FullAuto, &mode_id, &[])
6960 .expect("an empty mode catalog never offers anything");
6961 assert!(message.contains("ApprovalLevelNotOfferedByAgent"));
6962 assert!(message.contains("offered: []"));
6963 }
6964
6965 use super::{NativeEffectShell, NativeSessionKey, OwnedProviderSession};
6982 use gate4agent_catalog::{AgentRegistry as FixtureAgentRegistry, AgentSpec as FixtureAgentSpec};
6983 use gate4agent_types::ControlObservation;
6984 use std::collections::VecDeque;
6985
6986 fn fixture_pipe_key() -> NativeSessionKey {
6987 NativeSessionKey {
6988 instance_id: gate4agent_types::AgentInstanceId(9_001),
6989 generation: gate4agent_types::SessionGeneration(1),
6990 }
6991 }
6992
6993 fn fixture_pipe_source() -> gate4agent_types::ProviderSource {
6994 gate4agent_types::ProviderSource {
6995 family: AdapterFamily::Pipe,
6996 binding: gate4agent_types::AdapterBinding::new(
6997 gate4agent_types::AdapterId::new("fixture").unwrap(),
6998 "1",
6999 gate4agent_types::AdapterVerification::SyntheticFixture,
7000 )
7001 .unwrap(),
7002 }
7003 }
7004
7005 #[cfg(windows)]
7006 fn fixture_graceful_exit_launch() -> gate4agent_types::LaunchSpec {
7007 gate4agent_types::LaunchSpec {
7008 program: "powershell.exe".to_owned(),
7009 fixed_args: vec![
7010 "-NoProfile".to_owned(),
7011 "-NonInteractive".to_owned(),
7012 "-Command".to_owned(),
7013 "[Console]::In.ReadToEnd() | Out-Null; exit 0".to_owned(),
7014 ],
7015 }
7016 }
7017
7018 #[cfg(not(windows))]
7019 fn fixture_graceful_exit_launch() -> gate4agent_types::LaunchSpec {
7020 gate4agent_types::LaunchSpec {
7021 program: "sh".to_owned(),
7022 fixed_args: vec!["-c".to_owned(), "cat >/dev/null; exit 0".to_owned()],
7023 }
7024 }
7025
7026 #[cfg(windows)]
7027 fn fixture_ignore_stdin_launch() -> gate4agent_types::LaunchSpec {
7028 gate4agent_types::LaunchSpec {
7029 program: "powershell.exe".to_owned(),
7030 fixed_args: vec![
7031 "-NoProfile".to_owned(),
7032 "-NonInteractive".to_owned(),
7033 "-Command".to_owned(),
7034 "Start-Sleep -Seconds 30".to_owned(),
7035 ],
7036 }
7037 }
7038
7039 #[cfg(not(windows))]
7040 fn fixture_ignore_stdin_launch() -> gate4agent_types::LaunchSpec {
7041 gate4agent_types::LaunchSpec {
7042 program: "sh".to_owned(),
7043 fixed_args: vec!["-c".to_owned(), "sleep 30".to_owned()],
7044 }
7045 }
7046
7047 async fn spawn_fixture_pipe_shell(
7048 launch: gate4agent_types::LaunchSpec,
7049 ) -> (NativeEffectShell, NativeSessionKey) {
7050 let session = gate4agent::PipeSession::spawn_with_launch(
7051 gate4agent::core::types::SessionConfig::default(),
7052 "prompt",
7053 &launch,
7054 gate4agent_types::PipePromptDelivery::StdinClose,
7055 )
7056 .await
7057 .expect("fixture pipe process must spawn");
7058 let events = session.subscribe();
7059 let key = fixture_pipe_key();
7060 let mut shell = NativeEffectShell::new(
7061 FixtureAgentRegistry::new(std::iter::empty::<FixtureAgentSpec>())
7062 .expect("empty fixture catalog"),
7063 );
7064 shell.pipe_sessions.insert(
7065 key,
7066 OwnedProviderSession {
7067 source: fixture_pipe_source(),
7068 session,
7069 events,
7070 pending_events: VecDeque::new(),
7071 pending_provider_events: VecDeque::new(),
7072 next_provider_sequence: 1,
7073 observed_exit_code: None,
7074 runtime_policy: gate4agent_types::ProviderRuntimePolicy::none(),
7075 },
7076 );
7077 (shell, key)
7078 }
7079
7080 #[tokio::test]
7081 async fn stop_native_pipe_session_reports_the_real_exit_code_when_not_forced() {
7082 let (mut shell, key) = spawn_fixture_pipe_shell(fixture_graceful_exit_launch()).await;
7083
7084 match shell.stop_native(key, false).await {
7085 ControlObservation::StopCompleted {
7086 forced, exit_code, ..
7087 } => {
7088 assert!(!forced, "a process that exits on its own must not be forced");
7089 assert_eq!(exit_code, Some(0));
7090 }
7091 other => panic!("expected StopCompleted, got {other:?}"),
7092 }
7093 }
7094
7095 #[tokio::test]
7096 async fn stop_native_pipe_session_force_true_kills_immediately_without_waiting() {
7097 let (mut shell, key) = spawn_fixture_pipe_shell(fixture_ignore_stdin_launch()).await;
7098
7099 let start = std::time::Instant::now();
7100 match shell.stop_native(key, true).await {
7101 ControlObservation::StopCompleted { forced, .. } => {
7102 assert!(forced, "force=true must always report forced");
7103 }
7104 other => panic!("expected StopCompleted, got {other:?}"),
7105 }
7106 assert!(
7107 start.elapsed()
7108 < std::time::Duration::from_secs(gate4agent::pipe::PIPE_GRACEFUL_STOP_BOUND_SECS),
7109 "force=true must not wait out the graceful bound"
7110 );
7111 }
7112
7113 #[cfg(windows)]
7129 fn acp_fixture_launch() -> gate4agent_types::LaunchSpec {
7130 gate4agent_types::LaunchSpec {
7131 program: "powershell.exe".to_owned(),
7132 fixed_args: vec![
7133 "-NoLogo".to_owned(),
7134 "-NoProfile".to_owned(),
7135 "-NonInteractive".to_owned(),
7136 "-ExecutionPolicy".to_owned(),
7137 "Bypass".to_owned(),
7138 "-Command".to_owned(),
7139 ACP_HANDSHAKE_SCRIPT.to_owned(),
7140 ],
7141 }
7142 }
7143
7144 #[cfg(not(windows))]
7145 fn acp_fixture_launch() -> gate4agent_types::LaunchSpec {
7146 gate4agent_types::LaunchSpec {
7147 program: "python3".to_owned(),
7148 fixed_args: vec!["-u".to_owned(), "-c".to_owned(), ACP_HANDSHAKE_SCRIPT.to_owned()],
7149 }
7150 }
7151
7152 #[cfg(windows)]
7153 const ACP_HANDSHAKE_SCRIPT: &str = r#"[Console]::OutputEncoding=[Text.Encoding]::UTF8
7154function Write-JsonLine($value) { [Console]::WriteLine(($value | ConvertTo-Json -Compress -Depth 12)) }
7155$initialize = [Console]::ReadLine() | ConvertFrom-Json
7156Write-JsonLine @{jsonrpc='2.0';id=$initialize.id;result=@{}}
7157$newSession = [Console]::ReadLine() | ConvertFrom-Json
7158Write-JsonLine @{jsonrpc='2.0';id=$newSession.id;result=@{sessionId='fixture-acp-session'}}
7159while ($true) {
7160 $line = [Console]::ReadLine()
7161 if ($null -eq $line) { exit 0 }
7162}"#;
7163 #[cfg(not(windows))]
7164 const ACP_HANDSHAKE_SCRIPT: &str = r#"import json,sys
7165def read_message():
7166 line=sys.stdin.readline()
7167 if not line: return None
7168 return json.loads(line)
7169def write_message(message):
7170 print(json.dumps(message),flush=True)
7171
7172initialize=read_message()
7173write_message({'jsonrpc':'2.0','id':initialize.get('id'),'result':{}})
7174new_session=read_message()
7175write_message({'jsonrpc':'2.0','id':new_session.get('id'),'result':{'sessionId':'fixture-acp-session'}})
7176while True:
7177 msg=read_message()
7178 if msg is None:
7179 sys.exit(0)"#;
7180
7181 fn fixture_acp_key() -> NativeSessionKey {
7182 NativeSessionKey {
7183 instance_id: gate4agent_types::AgentInstanceId(9_002),
7184 generation: gate4agent_types::SessionGeneration(1),
7185 }
7186 }
7187
7188 fn fixture_acp_source() -> gate4agent_types::ProviderSource {
7189 gate4agent_types::ProviderSource {
7190 family: AdapterFamily::Acp,
7191 binding: gate4agent_types::AdapterBinding::new(
7192 gate4agent_types::AdapterId::new("fixture").unwrap(),
7193 "1",
7194 gate4agent_types::AdapterVerification::SyntheticFixture,
7195 )
7196 .unwrap(),
7197 }
7198 }
7199
7200 async fn spawn_fixture_acp_shell() -> (NativeEffectShell, NativeSessionKey) {
7201 let session = gate4agent::acp::AcpSession::spawn_with_launch(
7202 gate4agent::CliTool::ClaudeCode,
7203 &std::env::current_dir().expect("cwd"),
7204 gate4agent::acp::AcpSessionOptions::default(),
7205 &acp_fixture_launch(),
7206 )
7207 .await
7208 .expect("fixture ACP handshake must succeed");
7209 let events = session.subscribe();
7210 let key = fixture_acp_key();
7211 let mut shell = NativeEffectShell::new(
7212 FixtureAgentRegistry::new(std::iter::empty::<FixtureAgentSpec>())
7213 .expect("empty fixture catalog"),
7214 );
7215 shell.acp_sessions.insert(
7216 key,
7217 OwnedProviderSession {
7218 source: fixture_acp_source(),
7219 session,
7220 events,
7221 pending_events: VecDeque::new(),
7222 pending_provider_events: VecDeque::new(),
7223 next_provider_sequence: 1,
7224 observed_exit_code: None,
7225 runtime_policy: gate4agent_types::ProviderRuntimePolicy::none(),
7226 },
7227 );
7228 (shell, key)
7229 }
7230
7231 #[tokio::test]
7232 async fn stop_native_acp_session_reports_the_real_exit_code_when_not_forced() {
7233 let (mut shell, key) = spawn_fixture_acp_shell().await;
7234
7235 match shell.stop_native(key, false).await {
7236 ControlObservation::StopCompleted {
7237 forced, exit_code, ..
7238 } => {
7239 assert!(!forced, "a process that exits on its own must not be forced");
7240 assert_eq!(exit_code, Some(0));
7241 }
7242 other => panic!("expected StopCompleted, got {other:?}"),
7243 }
7244 }
7245
7246 #[cfg(windows)]
7264 fn acp_fixture_launch_no_session_id() -> gate4agent_types::LaunchSpec {
7265 gate4agent_types::LaunchSpec {
7266 program: "powershell.exe".to_owned(),
7267 fixed_args: vec![
7268 "-NoLogo".to_owned(),
7269 "-NoProfile".to_owned(),
7270 "-NonInteractive".to_owned(),
7271 "-ExecutionPolicy".to_owned(),
7272 "Bypass".to_owned(),
7273 "-Command".to_owned(),
7274 ACP_HANDSHAKE_SCRIPT_NO_SESSION_ID.to_owned(),
7275 ],
7276 }
7277 }
7278
7279 #[cfg(not(windows))]
7280 fn acp_fixture_launch_no_session_id() -> gate4agent_types::LaunchSpec {
7281 gate4agent_types::LaunchSpec {
7282 program: "python3".to_owned(),
7283 fixed_args: vec![
7284 "-u".to_owned(),
7285 "-c".to_owned(),
7286 ACP_HANDSHAKE_SCRIPT_NO_SESSION_ID.to_owned(),
7287 ],
7288 }
7289 }
7290
7291 #[cfg(windows)]
7292 const ACP_HANDSHAKE_SCRIPT_NO_SESSION_ID: &str = r#"[Console]::OutputEncoding=[Text.Encoding]::UTF8
7293function Write-JsonLine($value) { [Console]::WriteLine(($value | ConvertTo-Json -Compress -Depth 12)) }
7294$initialize = [Console]::ReadLine() | ConvertFrom-Json
7295Write-JsonLine @{jsonrpc='2.0';id=$initialize.id;result=@{}}
7296$newSession = [Console]::ReadLine() | ConvertFrom-Json
7297Write-JsonLine @{jsonrpc='2.0';id=$newSession.id;result=@{}}
7298while ($true) {
7299 $line = [Console]::ReadLine()
7300 if ($null -eq $line) { exit 0 }
7301}"#;
7302 #[cfg(not(windows))]
7303 const ACP_HANDSHAKE_SCRIPT_NO_SESSION_ID: &str = r#"import json,sys
7304def read_message():
7305 line=sys.stdin.readline()
7306 if not line: return None
7307 return json.loads(line)
7308def write_message(message):
7309 print(json.dumps(message),flush=True)
7310
7311initialize=read_message()
7312write_message({'jsonrpc':'2.0','id':initialize.get('id'),'result':{}})
7313new_session=read_message()
7314write_message({'jsonrpc':'2.0','id':new_session.get('id'),'result':{}})
7315while True:
7316 msg=read_message()
7317 if msg is None:
7318 sys.exit(0)"#;
7319
7320 fn fixture_acp_spawn_agent_id() -> AgentId {
7321 AgentId::new("claude").expect("builtin catalog names 'claude'")
7322 }
7323
7324 fn fixture_acp_spawn_agent_spec(launch: gate4agent_types::LaunchSpec) -> FixtureAgentSpec {
7330 let mut spec = gate4agent_catalog::builtin_registry()
7331 .get(&fixture_acp_spawn_agent_id())
7332 .cloned()
7333 .expect("builtin catalog must declare claude");
7334 spec.capabilities
7335 .transports
7336 .acp
7337 .as_mut()
7338 .expect("claude must declare ACP capability")
7339 .launch_override = Some(launch);
7340 spec
7341 }
7342
7343 fn fixture_acp_spawn_key() -> NativeSessionKey {
7344 NativeSessionKey {
7345 instance_id: gate4agent_types::AgentInstanceId(9_101),
7346 generation: gate4agent_types::SessionGeneration(1),
7347 }
7348 }
7349
7350 async fn spawn_native_acp_fixture(
7356 launch: gate4agent_types::LaunchSpec,
7357 provider_session_identity: bool,
7358 ) -> (NativeEffectShell, NativeSessionKey, ControlObservation) {
7359 let mut policy = ProviderRuntimePolicy::none();
7360 policy.provider_session_identity = provider_session_identity;
7361 let key = fixture_acp_spawn_key();
7362 let mut shell = NativeEffectShell::new(
7363 FixtureAgentRegistry::new([fixture_acp_spawn_agent_spec(launch)])
7364 .expect("single-agent fixture catalog"),
7365 );
7366 let envelope = gate4agent_types::EffectEnvelope {
7367 operation_id: gate4agent_types::OperationId(1),
7368 instance_id: key.instance_id,
7369 generation: key.generation,
7370 effect: gate4agent_types::ControlEffect::Spawn {
7371 agent_id: fixture_acp_spawn_agent_id(),
7372 transport: TransportKind::Acp,
7373 runtime_policy: policy,
7374 request: gate4agent_types::StartRequest {
7375 working_directory: std::env::current_dir()
7376 .expect("cwd")
7377 .to_string_lossy()
7378 .into_owned(),
7379 terminal_size: gate4agent_types::TerminalSize {
7380 rows: 24,
7381 columns: 80,
7382 },
7383 initial_prompt: None,
7384 session_options: None,
7385 approval_level: ApprovalLevel::Unmanaged,
7386 },
7387 },
7388 };
7389 let observation = shell.execute(envelope).await.observation;
7390 (shell, key, observation)
7391 }
7392
7393 #[tokio::test]
7394 async fn acp_spawn_with_a_real_sessionid_emits_session_identity_observed_after_session_start() {
7395 let (mut shell, key, spawned) =
7396 spawn_native_acp_fixture(acp_fixture_launch(), true).await;
7397 match spawned {
7398 ControlObservation::Spawned { .. } => {}
7399 other => panic!("expected Spawned, got {other:?}"),
7400 }
7401
7402 let mut provider_events = shell
7403 .collect_provider_events()
7404 .into_iter()
7405 .filter_map(|envelope| match envelope.observation {
7406 ControlObservation::ProviderEvent { event, .. } => Some(event),
7407 _ => None,
7408 });
7409
7410 match provider_events.next() {
7411 Some(ProviderEvent::SessionStarted { .. }) => {}
7412 other => panic!("expected SessionStarted first, got {other:?}"),
7413 }
7414 match provider_events.next() {
7415 Some(ProviderEvent::SessionIdentityObserved { identity }) => {
7416 assert_eq!(identity.key, gate4agent_types::ProviderSessionKey::SessionId);
7417 assert_eq!(identity.id, "fixture-acp-session");
7418 assert!(identity.transcript_path.is_none());
7419 }
7420 other => panic!("expected SessionIdentityObserved second, got {other:?}"),
7421 }
7422 assert!(
7423 provider_events.next().is_none(),
7424 "expected exactly one SessionIdentityObserved and nothing after it"
7425 );
7426
7427 let _ = shell.stop_native(key, true).await;
7428 }
7429
7430 #[tokio::test]
7431 async fn acp_spawn_with_no_sessionid_emits_no_session_identity_observed() {
7432 let (mut shell, key, spawned) =
7433 spawn_native_acp_fixture(acp_fixture_launch_no_session_id(), true).await;
7434 match spawned {
7435 ControlObservation::Spawned { .. } => {}
7436 other => panic!("expected Spawned, got {other:?}"),
7437 }
7438
7439 let provider_events: Vec<_> = shell
7440 .collect_provider_events()
7441 .into_iter()
7442 .filter_map(|envelope| match envelope.observation {
7443 ControlObservation::ProviderEvent { event, .. } => Some(event),
7444 _ => None,
7445 })
7446 .collect();
7447
7448 assert!(
7449 provider_events
7450 .iter()
7451 .any(|event| matches!(event, ProviderEvent::SessionStarted { .. })),
7452 "expected SessionStarted even when the agent reported no sessionId, got {provider_events:?}"
7453 );
7454 assert!(
7455 !provider_events
7456 .iter()
7457 .any(|event| matches!(event, ProviderEvent::SessionIdentityObserved { .. })),
7458 "an agent that violates ACP's sessionId MUST must not be granted an identity: \
7459 got {provider_events:?}"
7460 );
7461
7462 let _ = shell.stop_native(key, true).await;
7463 }
7464}