1use std::sync::mpsc::TryRecvError;
4use std::sync::{Arc, Mutex};
5use std::time::{Duration, Instant};
6
7use gate4agent_pty::ChildKiller;
8use tokio::sync::broadcast;
9use tokio::task::JoinHandle;
10
11use crate::agent::{
12 builtin_registry, plan_draft_launch, plan_launch, prepare_agent_command, prepare_input,
13 prepare_shell_command, AgentCommandMode, AgentId, AgentSpec, InputAction, LaunchPlan,
14 LaunchRequest, PreparedInput, PreparedInputKind, PromptFraming, PromptPayload, ReadinessIntent,
15 ReadinessPermit, ShellCommand, TerminalControl, TERMINAL_WRITE_DELAY_MAX_MS,
16};
17use crate::core::error::AgentError;
18use crate::core::types::{AgentEvent, CliTool, SessionConfig};
19use crate::pty::cli::traits::{MessageClass, StartupAction};
20use crate::pty::cli::{create_pipeline, create_submitter};
21use crate::pty::rate_limit::RateLimitDetector;
22use crate::pty::vte::VteParser;
23
24use super::event::{
25 PtyAttachError, PtyAttachment, PtyEvent, PtyEventPublisher, PtyEventReceiver, PtyReplayCursor,
26 PtySignal, PtySignalOutcome, PtySize, PtyTerminalSnapshot, DEFAULT_PTY_REPLAY_BYTES,
27 PTY_PROVIDER_PROTOCOL_REVISION,
28};
29use super::os_process::PtyForegroundObservation;
30use super::process_tree::{terminate_process_tree, PtyTreeTerminationReport};
31use super::wrapper::{PtyError, PtyReadEvent, PtyWrapper};
32
33const PTY_EVENT_CHANNEL_CAPACITY: usize = 64;
34const PTY_POST_EXIT_DRAIN_QUIET_MS: u64 = 100;
35pub const PTY_SHUTDOWN_TIMEOUT_MS: u64 = 8_000;
36
37#[derive(Clone, Debug, Eq, PartialEq)]
38pub struct PtyShutdownOutcome {
39 pub exit_code: Option<i32>,
40 pub termination: Option<PtyTreeTerminationReport>,
41 pub terminal: PtyTerminalSnapshot,
42}
43
44#[derive(Clone, Copy, Debug, Eq, PartialEq)]
53pub struct ForegroundProbeTiming {
54 pub queued: Duration,
58 pub lock_wait: Duration,
61 pub walk: Duration,
64}
65
66impl From<PtyError> for AgentError {
68 fn from(e: PtyError) -> Self {
69 match e {
70 PtyError::CreateFailed(s) => AgentError::PtyCreate(s),
71 PtyError::SpawnFailed(s) => AgentError::PtySpawn(s),
72 PtyError::Io(e) => AgentError::PtyIo { source: e },
73 PtyError::Pty(s) => AgentError::Pty(s),
74 PtyError::UnsafeWindowsCommandArgument { index } => AgentError::PtySpawn(format!(
75 "Windows command wrapper argument {index} contains shell metacharacters"
76 )),
77 PtyError::UnsupportedWindowsUncWorkingDirectory { path } => AgentError::PtySpawn(
78 format!("Windows PTY working directory cannot be a UNC path: {path}"),
79 ),
80 PtyError::ReaderJoinTimedOut { timeout_ms } => {
81 AgentError::PtyShutdownTimedOut { timeout_ms }
82 }
83 PtyError::ReaderPanicked => AgentError::Pty("PTY OS reader thread panicked".into()),
84 }
85 }
86}
87
88pub struct PtyWriteHandle {
90 inner: Arc<Mutex<PtyWrapper>>,
91}
92
93impl PtyWriteHandle {
94 pub fn write(&self, data: &str) -> Result<(), AgentError> {
96 let mut pty = self
97 .inner
98 .lock()
99 .map_err(|_| AgentError::Pty("PTY mutex poisoned".into()))?;
100 pty.write(data).map_err(|e| AgentError::Pty(e.to_string()))
101 }
102
103 pub fn write_bytes(&self, data: &[u8]) -> Result<(), AgentError> {
105 let mut pty = self
106 .inner
107 .lock()
108 .map_err(|_| AgentError::Pty("PTY mutex poisoned".into()))?;
109 pty.write_bytes(data).map_err(AgentError::from)
110 }
111
112 pub fn resize(&self, rows: u16, cols: u16) -> Result<(), AgentError> {
114 let pty = self
115 .inner
116 .lock()
117 .map_err(|_| AgentError::Pty("PTY mutex poisoned".into()))?;
118 pty.resize(rows, cols).map_err(AgentError::from)
119 }
120}
121
122pub struct PtySession {
136 session_id: String,
137 agent_id: AgentId,
138 process_spec: AgentSpec,
139 legacy_tool: Option<CliTool>,
140 agent_command_mode: Option<AgentCommandMode>,
141 tx: broadcast::Sender<AgentEvent>,
142 pty: Arc<Mutex<PtyWrapper>>,
143 root_pid: Option<u32>,
144 reader_task: Option<JoinHandle<()>>,
145 killer: Arc<Mutex<Box<dyn ChildKiller + Send + Sync>>>,
146 events: Arc<PtyEventPublisher>,
147 pending_followup_prompt: Option<String>,
148 pending_followup_prompt_inserted: bool,
149 pending_followup_draft: Option<String>,
150 typed_input_lock: tokio::sync::Mutex<()>,
151}
152
153impl PtySession {
154 pub async fn spawn(config: SessionConfig) -> Result<Self, AgentError> {
158 Self::spawn_with_size(config, 24, 80).await
159 }
160
161 pub async fn spawn_with_size(
163 config: SessionConfig,
164 rows: u16,
165 cols: u16,
166 ) -> Result<Self, AgentError> {
167 let tool = config.tool;
168 let agent_id = AgentId::from(tool);
169 let process_spec = builtin_registry()
170 .get(&agent_id)
171 .expect("legacy CLI tools have built-in process specifications")
172 .clone();
173 let pty =
174 PtyWrapper::new_with_env(tool, &config.working_dir, &config.env_vars, rows, cols)?;
175 Ok(Self::start(
176 agent_id,
177 process_spec,
178 Some(tool),
179 match tool {
180 CliTool::KimiCode | CliTool::Grok => None,
181 CliTool::ClaudeCode | CliTool::Codex => Some(AgentCommandMode::SlashLine),
182 },
183 pty,
184 None,
185 None,
186 rows,
187 cols,
188 ))
189 }
190
191 pub async fn spawn_agent(spec: &AgentSpec, request: LaunchRequest) -> Result<Self, AgentError> {
193 Self::spawn_agent_with_size(spec, request, 24, 80).await
194 }
195
196 pub async fn spawn_agent_with_size(
198 spec: &AgentSpec,
199 request: LaunchRequest,
200 rows: u16,
201 cols: u16,
202 ) -> Result<Self, AgentError> {
203 let plan = plan_launch(spec, request)?;
204 Self::spawn_generic_plan(
205 plan,
206 spec.clone(),
207 spec.capabilities.agent_commands,
208 rows,
209 cols,
210 )
211 }
212
213 pub async fn spawn_agent_draft(
215 spec: &AgentSpec,
216 request: LaunchRequest,
217 draft: String,
218 ) -> Result<Self, AgentError> {
219 Self::spawn_agent_draft_with_size(spec, request, draft, 24, 80).await
220 }
221
222 pub async fn spawn_agent_draft_with_size(
224 spec: &AgentSpec,
225 request: LaunchRequest,
226 draft: String,
227 rows: u16,
228 cols: u16,
229 ) -> Result<Self, AgentError> {
230 let plan = plan_draft_launch(spec, request, draft)?;
231 Self::spawn_generic_plan(
232 plan,
233 spec.clone(),
234 spec.capabilities.agent_commands,
235 rows,
236 cols,
237 )
238 }
239
240 fn spawn_generic_plan(
241 plan: LaunchPlan,
242 process_spec: AgentSpec,
243 agent_command_mode: Option<AgentCommandMode>,
244 rows: u16,
245 cols: u16,
246 ) -> Result<Self, AgentError> {
247 let agent_id = plan.agent_id.clone();
248 let pending_followup_prompt = plan.followup_prompt.clone();
249 let pending_followup_draft = plan.followup_draft.clone();
250 let legacy_tool = None;
254 let pty = PtyWrapper::from_launch_plan(plan, legacy_tool, rows, cols)?;
255 Ok(Self::start(
256 agent_id,
257 process_spec,
258 legacy_tool,
259 agent_command_mode,
260 pty,
261 pending_followup_prompt,
262 pending_followup_draft,
263 rows,
264 cols,
265 ))
266 }
267
268 fn start(
269 agent_id: AgentId,
270 process_spec: AgentSpec,
271 legacy_tool: Option<CliTool>,
272 agent_command_mode: Option<AgentCommandMode>,
273 pty: PtyWrapper,
274 pending_followup_prompt: Option<String>,
275 pending_followup_draft: Option<String>,
276 rows: u16,
277 cols: u16,
278 ) -> Self {
279 let session_id = uuid_v4();
280 let provider_revision = format!(
281 "{PTY_PROVIDER_PROTOCOL_REVISION}:{}:{}",
282 agent_id.as_str(),
283 process_spec.revision
284 );
285 let root_pid = pty.root_pid();
286 let killer = Arc::new(Mutex::new(pty.clone_killer()));
287 let pty = Arc::new(Mutex::new(pty));
288 let (tx, _) = broadcast::channel::<AgentEvent>(4096);
289 let events = PtyEventPublisher::new(
290 session_id.clone(),
291 provider_revision,
292 1,
293 PTY_EVENT_CHANNEL_CAPACITY,
294 DEFAULT_PTY_REPLAY_BYTES,
295 rows,
296 cols,
297 );
298
299 let _ = tx.send(AgentEvent::Started {
300 session_id: session_id.clone(),
301 });
302 events.publish(PtyEvent::Started);
303
304 let pty_clone = pty.clone();
305 let tx_clone = tx.clone();
306 let events_clone = events.clone();
307 let reader_task = tokio::task::spawn_blocking(move || {
308 reader_loop(pty_clone, tx_clone, events_clone, legacy_tool);
309 });
310
311 Self {
312 session_id,
313 agent_id,
314 process_spec,
315 legacy_tool,
316 agent_command_mode,
317 tx,
318 pty,
319 root_pid,
320 reader_task: Some(reader_task),
321 killer,
322 events,
323 pending_followup_prompt,
324 pending_followup_prompt_inserted: false,
325 pending_followup_draft,
326 typed_input_lock: tokio::sync::Mutex::new(()),
327 }
328 }
329
330 pub fn subscribe(&self) -> broadcast::Receiver<AgentEvent> {
334 self.tx.subscribe()
335 }
336
337 pub fn subscribe_events(&self) -> Result<PtyEventReceiver, PtyAttachError> {
339 self.events.subscribe()
340 }
341
342 pub fn attach_events(&self, cursor: PtyReplayCursor) -> Result<PtyAttachment, PtyAttachError> {
344 self.events.attach(cursor)
345 }
346
347 pub fn terminal_snapshot(&self) -> Result<PtyTerminalSnapshot, PtyAttachError> {
349 let snapshot = self.events.snapshot()?;
350 self.events.publish(PtyEvent::SnapshotAvailable {
351 snapshot_sequence: snapshot.sequence,
352 });
353 Ok(snapshot)
354 }
355
356 pub fn terminal_state(&self) -> Result<PtyTerminalSnapshot, PtyAttachError> {
359 self.events.snapshot()
360 }
361
362 pub fn terminal_sequence(&self) -> Result<u64, PtyAttachError> {
367 self.events.terminal_sequence()
368 }
369
370 pub async fn observe_foreground(&self) -> Result<PtyForegroundObservation, AgentError> {
372 self.observe_foreground_timed()
373 .await
374 .map(|(observation, _timing)| observation)
375 }
376
377 pub async fn observe_foreground_timed(
389 &self,
390 ) -> Result<(PtyForegroundObservation, ForegroundProbeTiming), AgentError> {
391 let pty = self.pty.clone();
392 let spec = self.process_spec.clone();
393 let dispatched_at = Instant::now();
394 let (observation, timing) = tokio::task::spawn_blocking(move || {
395 let started_at = Instant::now();
396 let queued = started_at.duration_since(dispatched_at);
397 let lock_started_at = Instant::now();
398 let guard = pty
399 .lock()
400 .map_err(|_| AgentError::Pty("PTY mutex poisoned".into()))?;
401 let lock_wait = lock_started_at.elapsed();
402 let walk_started_at = Instant::now();
403 let observation = guard.observe_foreground(&spec).map_err(AgentError::from)?;
404 let walk = walk_started_at.elapsed();
405 Ok::<_, AgentError>((
406 observation,
407 ForegroundProbeTiming {
408 queued,
409 lock_wait,
410 walk,
411 },
412 ))
413 })
414 .await
415 .map_err(|_| AgentError::Pty("spawn_blocking panicked".into()))??;
416 self.events
417 .publish(PtyEvent::ForegroundProcess(observation.clone()));
418 Ok((observation, timing))
419 }
420
421 pub fn write_handle(&self) -> PtyWriteHandle {
423 PtyWriteHandle {
424 inner: self.pty.clone(),
425 }
426 }
427
428 pub async fn write(&self, data: &str) -> Result<(), AgentError> {
430 let data = data.to_owned();
431 let pty = self.pty.clone();
432 tokio::task::spawn_blocking(move || {
433 let mut guard = pty
434 .lock()
435 .map_err(|_| AgentError::Pty("PTY mutex poisoned".into()))?;
436 guard
437 .write(&data)
438 .map_err(|e| AgentError::Pty(e.to_string()))
439 })
440 .await
441 .map_err(|_| AgentError::Pty("spawn_blocking panicked".into()))?
442 }
443
444 pub async fn send_input_action(
446 &self,
447 action: InputAction,
448 permit: ReadinessPermit,
449 ) -> Result<(), AgentError> {
450 let prepared = match action {
451 InputAction::AgentCommand(command) => {
452 self.validate_permit(&permit, Some(ReadinessIntent::DraftPaste))?;
453 if self.agent_command_mode != Some(AgentCommandMode::SlashLine) {
454 return Err(AgentError::AgentCapabilityUnsupported {
455 agent: self.agent_id.clone(),
456 capability: "agent-commands",
457 });
458 }
459 prepare_agent_command(command, &self.agent_id)?
460 }
461 action => prepare_input(action)?,
462 };
463 self.send_prepared_input(prepared, permit).await
464 }
465
466 pub async fn send_prepared_input(
471 &self,
472 input: PreparedInput,
473 permit: ReadinessPermit,
474 ) -> Result<(), AgentError> {
475 let required_intent = match input.kind() {
476 PreparedInputKind::InsertDraft | PreparedInputKind::AgentCommand => {
477 Some(ReadinessIntent::DraftPaste)
478 }
479 PreparedInputKind::SubmitPrompt => Some(ReadinessIntent::FollowupPrompt),
480 PreparedInputKind::ShellCommand => {
481 return Err(AgentError::Pty(
482 "shell input requires fresh foreground-shell proof".to_owned(),
483 ));
484 }
485 PreparedInputKind::TerminalText
486 | PreparedInputKind::TerminalBytes
487 | PreparedInputKind::TerminalControl => None,
488 };
489 self.validate_permit(&permit, required_intent)?;
490 self.write_prepared_input(input).await
491 }
492
493 pub async fn send_terminal_input(&self, input: PreparedInput) -> Result<(), AgentError> {
496 if !matches!(
497 input.kind(),
498 PreparedInputKind::TerminalText
499 | PreparedInputKind::TerminalBytes
500 | PreparedInputKind::TerminalControl
501 ) {
502 return Err(AgentError::Pty(
503 "semantic input requires a readiness permit".to_owned(),
504 ));
505 }
506 self.write_prepared_input(input).await
507 }
508
509 pub async fn send_shell_input(&self, input: PreparedInput) -> Result<(), AgentError> {
512 if input.kind() != PreparedInputKind::ShellCommand {
513 return Err(AgentError::Pty(
514 "foreground-shell dispatch requires a prepared shell command".to_owned(),
515 ));
516 }
517 let _serialized = self.typed_input_lock.lock().await;
518 let foreground = self.observe_foreground().await?;
519 if !foreground.readiness.is_shell {
520 return Err(AgentError::Pty(format!(
521 "refusing shell command: PTY foreground '{}' is not a shell",
522 foreground.observed_process
523 )));
524 }
525 self.write_prepared_input_locked(input).await
526 }
527
528 pub async fn send_agent_command_input(
531 &self,
532 input: PreparedInput,
533 permit: ReadinessPermit,
534 ) -> Result<(), AgentError> {
535 if input.kind() != PreparedInputKind::AgentCommand {
536 return Err(AgentError::Pty(
537 "agent-command dispatch requires a prepared agent command".to_owned(),
538 ));
539 }
540 self.validate_permit(&permit, Some(ReadinessIntent::DraftPaste))?;
541 let _serialized = self.typed_input_lock.lock().await;
542 let foreground = self.observe_foreground().await?;
543 if foreground.readiness.process_name.as_deref() != Some(self.agent_id.as_str()) {
544 return Err(AgentError::Pty(format!(
545 "refusing agent command: PTY foreground '{}' is not agent '{}'",
546 foreground.observed_process, self.agent_id
547 )));
548 }
549 self.write_prepared_input_locked(input).await
550 }
551
552 pub async fn send_shell_command(&self, command: ShellCommand) -> Result<(), AgentError> {
555 self.send_shell_input(prepare_shell_command(command)?).await
556 }
557
558 async fn write_prepared_input(&self, input: PreparedInput) -> Result<(), AgentError> {
559 let _serialized = self.typed_input_lock.lock().await;
560 self.write_prepared_input_locked(input).await
561 }
562
563 async fn write_prepared_input_locked(&self, input: PreparedInput) -> Result<(), AgentError> {
564 for write in input.into_writes() {
565 if write.delay_before_ms > TERMINAL_WRITE_DELAY_MAX_MS {
566 return Err(AgentError::Pty(format!(
567 "prepared write delay {}ms exceeds {}ms",
568 write.delay_before_ms, TERMINAL_WRITE_DELAY_MAX_MS
569 )));
570 }
571 if write.delay_before_ms > 0 {
572 tokio::time::sleep(Duration::from_millis(write.delay_before_ms)).await;
573 }
574 let pty = self.pty.clone();
575 tokio::task::spawn_blocking(move || {
576 let mut guard = pty
577 .lock()
578 .map_err(|_| AgentError::Pty("PTY mutex poisoned".into()))?;
579 guard.write_bytes(&write.bytes).map_err(AgentError::from)
580 })
581 .await
582 .map_err(|_| AgentError::Pty("spawn_blocking panicked".into()))??;
583 }
584 Ok(())
585 }
586
587 pub fn pending_followup_prompt(&self) -> Option<&str> {
589 self.pending_followup_prompt.as_deref()
590 }
591
592 pub async fn insert_pending_followup_prompt(
599 &mut self,
600 framing: PromptFraming,
601 permit: &ReadinessPermit,
602 ) -> Result<bool, AgentError> {
603 self.validate_permit(permit, Some(ReadinessIntent::FollowupPrompt))?;
604 let Some(prompt) = self.pending_followup_prompt.clone() else {
605 return Ok(false);
606 };
607 if self.pending_followup_prompt_inserted {
608 return Ok(false);
609 }
610 let input = prepare_input(InputAction::InsertDraft(PromptPayload {
611 text: prompt,
612 framing,
613 }))?;
614 self.write_prepared_input(input).await?;
615 self.pending_followup_prompt_inserted = true;
616 Ok(true)
617 }
618
619 pub fn pending_followup_draft(&self) -> Option<&str> {
621 self.pending_followup_draft.as_deref()
622 }
623
624 pub async fn insert_pending_followup_draft(
626 &mut self,
627 framing: PromptFraming,
628 permit: ReadinessPermit,
629 ) -> Result<bool, AgentError> {
630 self.validate_permit(&permit, Some(ReadinessIntent::DraftPaste))?;
631 let Some(draft) = self.pending_followup_draft.clone() else {
632 return Ok(false);
633 };
634 self.send_input_action(
635 InputAction::InsertDraft(PromptPayload {
636 text: draft,
637 framing,
638 }),
639 permit,
640 )
641 .await?;
642 self.pending_followup_draft = None;
643 Ok(true)
644 }
645
646 pub async fn submit_pending_followup(
648 &mut self,
649 framing: PromptFraming,
650 permit: ReadinessPermit,
651 ) -> Result<bool, AgentError> {
652 self.validate_permit(&permit, Some(ReadinessIntent::FollowupPrompt))?;
653 let Some(prompt) = self.pending_followup_prompt.clone() else {
654 return Ok(false);
655 };
656 if self.pending_followup_prompt_inserted {
657 let input = prepare_input(InputAction::TerminalControl(TerminalControl::Enter))?;
658 self.write_prepared_input(input).await?;
659 } else {
660 self.send_input_action(
661 InputAction::SubmitPrompt(PromptPayload {
662 text: prompt,
663 framing,
664 }),
665 permit,
666 )
667 .await?;
668 }
669 self.pending_followup_prompt = None;
670 self.pending_followup_prompt_inserted = false;
671 Ok(true)
672 }
673
674 fn validate_permit(
675 &self,
676 permit: &ReadinessPermit,
677 required_intent: Option<ReadinessIntent>,
678 ) -> Result<(), AgentError> {
679 if permit.agent_id() != &self.agent_id {
680 return Err(AgentError::PtyReadinessAgentMismatch {
681 session_agent: self.agent_id.clone(),
682 permit_agent: permit.agent_id().clone(),
683 });
684 }
685 if let Some(required) = required_intent {
686 if permit.intent() != required {
687 return Err(AgentError::PtyReadinessIntentMismatch {
688 required,
689 actual: permit.intent(),
690 });
691 }
692 }
693 Ok(())
694 }
695
696 pub async fn send_prompt(&self, prompt: &str) -> Result<(), AgentError> {
701 if self.legacy_tool.is_none() {
702 return Err(AgentError::Pty(
703 "legacy send_prompt is unavailable for generic agent sessions; use typed readiness-gated input"
704 .into(),
705 ));
706 }
707 for ch in prompt.chars() {
708 let s = ch.to_string();
709 self.write(&s).await?;
710 tokio::time::sleep(Duration::from_millis(30)).await;
711 }
712 self.write("\r").await
713 }
714
715 pub async fn resize(&self, rows: u16, cols: u16) -> Result<(), AgentError> {
717 let pty = self.pty.clone();
718 tokio::task::spawn_blocking(move || {
719 let guard = pty
720 .lock()
721 .map_err(|_| AgentError::Pty("PTY mutex poisoned".into()))?;
722 guard.resize(rows, cols).map_err(AgentError::from)
723 })
724 .await
725 .map_err(|_| AgentError::Pty("spawn_blocking panicked".into()))??;
726 self.events
727 .publish(PtyEvent::Resized(PtySize { rows, cols }));
728 Ok(())
729 }
730
731 pub async fn signal(&self, signal: PtySignal) -> Result<PtySignalOutcome, AgentError> {
733 match signal {
734 PtySignal::InterruptKey => {
735 let pty = self.pty.clone();
736 tokio::task::spawn_blocking(move || {
737 let mut guard = pty
738 .lock()
739 .map_err(|_| AgentError::Pty("PTY mutex poisoned".into()))?;
740 guard.write_bytes(b"\x03").map_err(AgentError::from)
741 })
742 .await
743 .map_err(|_| AgentError::Pty("spawn_blocking panicked".into()))??;
744 Ok(PtySignalOutcome::ControlWritten)
745 }
746 PtySignal::EndOfFileKey => {
747 let pty = self.pty.clone();
748 tokio::task::spawn_blocking(move || {
749 let mut guard = pty
750 .lock()
751 .map_err(|_| AgentError::Pty("PTY mutex poisoned".into()))?;
752 guard.write_bytes(b"\x04").map_err(AgentError::from)
753 })
754 .await
755 .map_err(|_| AgentError::Pty("spawn_blocking panicked".into()))??;
756 Ok(PtySignalOutcome::ControlWritten)
757 }
758 PtySignal::TerminateProcess => {
759 self.kill().await?;
760 Ok(PtySignalOutcome::TerminationRequested)
761 }
762 }
763 }
764
765 pub fn session_id(&self) -> &str {
767 &self.session_id
768 }
769
770 pub fn generation(&self) -> u64 {
772 self.events.generation()
773 }
774
775 pub fn provider_revision(&self) -> &str {
777 self.events.provider_revision()
778 }
779
780 pub fn beginning_cursor(&self) -> PtyReplayCursor {
782 PtyReplayCursor::beginning(self.provider_revision(), self.generation())
783 }
784
785 pub fn retained_cursor(&self) -> Result<PtyReplayCursor, PtyAttachError> {
789 self.events.retained_cursor()
790 }
791
792 pub fn attach_retained_events(&self) -> Result<PtyAttachment, PtyAttachError> {
795 self.events.attach_retained()
796 }
797
798 pub fn agent_id(&self) -> &AgentId {
800 &self.agent_id
801 }
802
803 pub fn root_pid(&self) -> Option<u32> {
805 self.root_pid
806 }
807
808 pub fn reader_finished(&self) -> bool {
810 self.reader_task
811 .as_ref()
812 .is_none_or(JoinHandle::is_finished)
813 }
814
815 pub async fn terminate_tree(&self) -> Result<PtyTreeTerminationReport, AgentError> {
818 let root_pid = self.root_pid;
819 let killer = self.killer.clone();
820 tokio::task::spawn_blocking(move || {
821 terminate_process_tree(root_pid, || {
822 let mut killer = killer
823 .lock()
824 .map_err(|_| "PTY killer mutex poisoned".to_owned())?;
825 match killer.kill() {
826 Ok(()) => Ok(()),
827 #[cfg(windows)]
828 Err(error) if error.raw_os_error() == Some(0) => Ok(()),
829 Err(error) => Err(error.to_string()),
830 }
831 })
832 .map_err(AgentError::from)
833 })
834 .await
835 .map_err(|_| AgentError::Pty("spawn_blocking panicked".into()))?
836 }
837
838 pub async fn kill(&self) -> Result<(), AgentError> {
841 self.terminate_tree().await.map(|_| ())
842 }
843
844 pub async fn shutdown(mut self) -> Result<PtyShutdownOutcome, AgentError> {
847 let shutdown_started = Instant::now();
848 let root_pid = self.root_pid;
849 let termination = if self.reader_finished() {
850 None
851 } else {
852 Some(self.terminate_tree().await?)
853 };
854 if let Some(mut reader_task) = self.reader_task.take() {
855 tokio::time::timeout(
856 Duration::from_millis(PTY_SHUTDOWN_TIMEOUT_MS),
857 &mut reader_task,
858 )
859 .await
860 .map_err(|_| {
861 eprintln!(
866 "[gate4agent-pty-session] shutdown timed out joining the reader thread for \
867 root_pid={root_pid:?} after {PTY_SHUTDOWN_TIMEOUT_MS}ms (termination={termination:?})",
868 );
869 AgentError::PtyShutdownTimedOut {
870 timeout_ms: PTY_SHUTDOWN_TIMEOUT_MS,
871 }
872 })?
873 .map_err(|_| AgentError::Pty("PTY reader task panicked".into()))?;
874 }
875 let remaining = Duration::from_millis(PTY_SHUTDOWN_TIMEOUT_MS)
876 .saturating_sub(shutdown_started.elapsed());
877 let pty = self.pty.clone();
878 let exit_code = tokio::task::spawn_blocking(move || {
879 let mut guard = pty
880 .lock()
881 .map_err(|_| AgentError::Pty("PTY mutex poisoned".into()))?;
882 guard.close_and_join_reader(remaining)?;
883 Ok::<_, AgentError>(guard.try_exit_code().map(|code| code as i32))
884 })
885 .await
886 .map_err(|_| AgentError::Pty("spawn_blocking panicked".into()))??;
887 let terminal = self
888 .events
889 .snapshot()
890 .map_err(|error| AgentError::Pty(error.to_string()))?;
891 Ok(PtyShutdownOutcome {
892 exit_code,
893 termination,
894 terminal,
895 })
896 }
897}
898
899impl Drop for PtySession {
900 fn drop(&mut self) {
901 if self
902 .reader_task
903 .as_ref()
904 .is_none_or(JoinHandle::is_finished)
905 {
906 return;
907 }
908 let mut killer = match self.killer.lock() {
909 Ok(killer) => killer,
910 Err(poisoned) => poisoned.into_inner(),
911 };
912 let _ = killer.kill();
913 }
914}
915
916fn reader_loop(
921 pty: Arc<Mutex<PtyWrapper>>,
922 tx: broadcast::Sender<AgentEvent>,
923 events: Arc<PtyEventPublisher>,
924 legacy_tool: Option<CliTool>,
925) {
926 let mut vte_parser = VteParser::new();
927 let rate_limit_detector = legacy_tool.map(RateLimitDetector::new_for_tool);
928 let mut pipeline = legacy_tool.map(create_pipeline);
929 let submitter = legacy_tool.map(create_submitter);
930 let mut startup_done = false;
931 let mut startup_input_suppressed = false;
932 let mut output_closed = false;
933 let mut observed_exit: Option<(u32, Instant)> = None;
934
935 loop {
936 if output_closed {
937 let exit_code = pty.lock().ok().and_then(|mut guard| guard.try_exit_code());
938 if let Some(code) = exit_code {
939 let code = code as i32;
940 let _ = tx.send(AgentEvent::Exited { code });
941 events.publish(PtyEvent::Exited { code });
942 break;
943 }
944 std::thread::sleep(Duration::from_millis(10));
945 continue;
946 }
947
948 let exit_code = pty.lock().ok().and_then(|mut guard| guard.try_exit_code());
949 if let Some(code) = exit_code {
950 let (_, quiet_since) = observed_exit.get_or_insert((code, Instant::now()));
951 if quiet_since.elapsed() >= Duration::from_millis(PTY_POST_EXIT_DRAIN_QUIET_MS) {
952 let queue_closed = pty
953 .lock()
954 .ok()
955 .is_some_and(|guard| guard.close_output_if_empty());
956 if queue_closed {
957 let code = code as i32;
958 let _ = tx.send(AgentEvent::Exited { code });
959 events.publish(PtyEvent::Exited { code });
960 break;
961 }
962 }
963 }
964
965 let receive = match pty.lock() {
966 Ok(guard) => guard.try_recv_result(),
967 Err(_) => break,
968 };
969 let raw = match receive {
970 Ok(PtyReadEvent::Output(raw)) => {
971 if let Some((_, quiet_since)) = observed_exit.as_mut() {
972 *quiet_since = Instant::now();
973 }
974 raw
975 }
976 Ok(PtyReadEvent::Eof) => {
977 output_closed = true;
978 continue;
979 }
980 Ok(PtyReadEvent::Error(message)) => {
981 let _ = tx.send(AgentEvent::Error {
982 message: format!("PTY reader failed: {message}"),
983 });
984 events.publish(PtyEvent::ReaderError { message });
985 if let Ok(mut guard) = pty.lock() {
986 let _ = guard.kill();
987 }
988 output_closed = true;
989 continue;
990 }
991 Err(TryRecvError::Empty) => {
992 std::thread::sleep(Duration::from_millis(10));
993 continue;
994 }
995 Err(TryRecvError::Disconnected) => {
996 let message = "PTY reader channel closed without a terminal event".to_owned();
997 let _ = tx.send(AgentEvent::Error {
998 message: message.clone(),
999 });
1000 events.publish(PtyEvent::ReaderError { message });
1001 if let Ok(mut guard) = pty.lock() {
1002 let _ = guard.kill();
1003 }
1004 output_closed = true;
1005 continue;
1006 }
1007 };
1008
1009 let _ = tx.send(AgentEvent::PtyRaw { data: raw.clone() });
1013 events.publish(PtyEvent::Output(raw.clone()));
1014
1015 let Some(pipeline) = pipeline.as_mut() else {
1016 continue;
1017 };
1018
1019 let raw_str = String::from_utf8_lossy(&raw).to_string();
1023
1024 let cleaned = vte_parser.parse(&raw_str);
1026
1027 if let Some(rl_info) = rate_limit_detector
1029 .as_ref()
1030 .and_then(|detector| detector.detect(&cleaned))
1031 {
1032 let _ = tx.send(AgentEvent::RateLimit(rl_info));
1033 }
1034
1035 let messages = pipeline.process(&raw_str);
1037 for msg in messages {
1038 match msg.class {
1039 MessageClass::PromptReady => {
1040 let _ = tx.send(AgentEvent::PtyReady);
1041 }
1042 MessageClass::ToolApproval => {
1043 let tool_name = msg
1044 .metadata
1045 .tool_name
1046 .clone()
1047 .unwrap_or_else(|| "unknown".into());
1048 let _ = tx.send(AgentEvent::PtyToolApproval {
1049 tool_name,
1050 description: None,
1051 });
1052 }
1053 _ => {}
1054 }
1055 let _ = tx.send(AgentEvent::PtyParsed(msg));
1056 }
1057
1058 if !startup_done {
1060 let Some(submitter) = submitter.as_ref() else {
1061 continue;
1062 };
1063 let action = submitter.handle_startup(&cleaned);
1064 match action {
1065 StartupAction::Ready => {
1066 startup_done = true;
1067 }
1068 StartupAction::SendInput(_) => {
1069 if !startup_input_suppressed {
1070 startup_input_suppressed = true;
1071 let message =
1072 "automatic startup input was suppressed; operator action is required"
1073 .to_owned();
1074 let _ = tx.send(AgentEvent::Error {
1075 message: message.clone(),
1076 });
1077 events.publish(PtyEvent::OperatorActionRequired { message });
1078 }
1079 }
1080 StartupAction::Waiting => {}
1081 }
1082 }
1083 }
1084}
1085
1086fn uuid_v4() -> String {
1088 use std::time::{SystemTime, UNIX_EPOCH};
1089 let t = SystemTime::now()
1090 .duration_since(UNIX_EPOCH)
1091 .unwrap_or_default()
1092 .as_nanos();
1093 format!("pty-{:x}", t)
1094}