Skip to main content

gate4agent/pty/
session.rs

1//! Async PTY session with tokio broadcast fan-out.
2
3use 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/// Where the wall-clock time in one foreground probe actually went.
45///
46/// A single total for `observe_foreground` mixes three things that lead to
47/// opposite conclusions about the same measurement: half a second of
48/// `queued` is idle waiting behind unrelated work on tokio's blocking pool,
49/// while half a second of `walk` is this probe genuinely burning CPU on the
50/// OS process-table walk. `walk` is the only field of the three that is
51/// unambiguously CPU spent by this probe.
52#[derive(Clone, Copy, Debug, Eq, PartialEq)]
53pub struct ForegroundProbeTiming {
54    /// Call until the blocking closure began running: tokio's blocking-pool
55    /// queue. Waiting, not work -- a PTY read parked on a pool thread can
56    /// put a probe behind it without either of them burning a cycle.
57    pub queued: Duration,
58    /// Inside the closure, waiting for the PTY mutex, which live PTY I/O
59    /// also holds.
60    pub lock_wait: Duration,
61    /// The OS process-table walk itself. The only one of the three that is
62    /// unambiguously CPU this probe spent.
63    pub walk: Duration,
64}
65
66/// Conversion from the internal PtyError to the public AgentError.
67impl 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
88/// Opaque write handle for sending input to a PTY.
89pub struct PtyWriteHandle {
90    inner: Arc<Mutex<PtyWrapper>>,
91}
92
93impl PtyWriteHandle {
94    /// Write raw bytes to the PTY.
95    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    /// Write raw bytes to the PTY.
104    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    /// Resize the PTY.
113    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
122/// Async PTY session. Spawns a CLI tool in a real PTY and broadcasts `AgentEvent`
123/// to all subscribers via a tokio broadcast channel.
124///
125/// The reader loop runs on a blocking thread (via `spawn_blocking`) because PTY
126/// I/O is inherently blocking. The loop bridges to async consumers via the
127/// broadcast channel.
128///
129/// # Broadcast semantics
130///
131/// `subscribe()` preserves the legacy best-effort `AgentEvent` channel.
132/// Resilient consumers use `subscribe_events()` or `attach_events()`: those
133/// streams expose generation/sequence, convert lag into `DataGap`, and provide
134/// bounded replay plus a sequence-pinned terminal snapshot.
135pub 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    /// Spawn a CLI tool in a PTY and start broadcasting events.
155    ///
156    /// Uses a 24x80 terminal size (compact). To control size, use `spawn_with_size`.
157    pub async fn spawn(config: SessionConfig) -> Result<Self, AgentError> {
158        Self::spawn_with_size(config, 24, 80).await
159    }
160
161    /// Spawn a CLI tool in a PTY with a specific terminal size.
162    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    /// Spawn any registered agent from its shell-free launch specification.
192    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    /// Spawn any registered agent with a specific PTY size.
197    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    /// Spawn an agent with a reviewable draft that is never auto-submitted.
214    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    /// Spawn an agent with a reviewable draft and a specific PTY size.
223    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        // Generic registry entries never acquire a semantic adapter merely by
251        // reusing a built-in ID. Adapter registration is a separate capability
252        // boundary; the legacy four-tool constructor retains existing behavior.
253        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    /// Subscribe to receive all future `AgentEvent` values from this session.
331    ///
332    /// Note: events that occurred before subscribing will not be received.
333    pub fn subscribe(&self) -> broadcast::Receiver<AgentEvent> {
334        self.tx.subscribe()
335    }
336
337    /// Subscribe to sequenced PTY runtime events from the next event onward.
338    pub fn subscribe_events(&self) -> Result<PtyEventReceiver, PtyAttachError> {
339        self.events.subscribe()
340    }
341
342    /// Atomically obtain bounded replay and subscribe after the replay boundary.
343    pub fn attach_events(&self, cursor: PtyReplayCursor) -> Result<PtyAttachment, PtyAttachError> {
344        self.events.attach(cursor)
345    }
346
347    /// Capture terminal state and the exact event sequence it incorporates.
348    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    /// Capture terminal state without publishing a snapshot-available event.
357    /// Native runtimes use this to refresh replaceable control-plane state.
358    pub fn terminal_state(&self) -> Result<PtyTerminalSnapshot, PtyAttachError> {
359        self.events.snapshot()
360    }
361
362    /// The sequence `terminal_state` would report, without paying for the
363    /// capture. A caller that only wants to know whether anything has
364    /// happened since it last looked must ask this first -- see
365    /// `PtyEventPublisher::terminal_sequence`.
366    pub fn terminal_sequence(&self) -> Result<u64, PtyAttachError> {
367        self.events.terminal_sequence()
368    }
369
370    /// Take a fresh, bounded OS process-table observation for readiness.
371    pub async fn observe_foreground(&self) -> Result<PtyForegroundObservation, AgentError> {
372        self.observe_foreground_timed()
373            .await
374            .map(|(observation, _timing)| observation)
375    }
376
377    /// The same probe as [`Self::observe_foreground`], with the wall-clock
378    /// time it spent broken into where it actually went. A single total for
379    /// this probe is not interpretable: `spawn_blocking` queueing is tokio's
380    /// blocking pool -- pure waiting, possibly behind a parked PTY read, that
381    /// can cost no CPU at all -- and `lock_wait` is contention with live PTY
382    /// I/O for the same mutex, also waiting. Only `walk`, the OS
383    /// process-table scan itself, is unambiguously CPU this probe spent.
384    ///
385    /// Implemented as the body that `observe_foreground` delegates to, so
386    /// the two cannot drift apart and the same `ForegroundProcess` event is
387    /// published exactly once regardless of which one is called.
388    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    /// Get the write handle for sending input to the PTY.
422    pub fn write_handle(&self) -> PtyWriteHandle {
423        PtyWriteHandle {
424            inner: self.pty.clone(),
425        }
426    }
427
428    /// Send raw bytes to the PTY.
429    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    /// Prepare and write a typed terminal action after readiness is proven.
445    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    /// Write an already prepared operation after readiness is proven.
467    ///
468    /// Typed operations are serialized. Submit remains a distinct delayed write
469    /// so TUI paste handlers cannot consume Enter as part of the body.
470    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    /// Write explicit terminal text or control without an agent readiness
494    /// permit. Semantic prompt, draft, and agent-command inputs remain gated.
495    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    /// Confirm that a shell currently owns this PTY, then write one bounded
510    /// command while holding the typed-input serialization guard.
511    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    /// Reconfirm that this session's configured agent still owns the PTY
529    /// before sending a provider-native inline command.
530    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    /// Prepare and send an intentional shell command through the same fresh
553    /// foreground-proof path used by canonical effect executors.
554    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    /// Prompt deferred by the launch plan, if one has not been submitted yet.
588    pub fn pending_followup_prompt(&self) -> Option<&str> {
589        self.pending_followup_prompt.as_deref()
590    }
591
592    /// Insert the launch plan's deferred prompt without submitting it.
593    ///
594    /// Startup orchestration can use the resulting terminal render as proof
595    /// that a TUI consumed the paste before sending Enter. Repeated calls do
596    /// not paste the prompt twice and return `false`. This is a startup-only
597    /// operation; callers must discard the session after any write error.
598    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    /// Reviewable draft deferred by the launch plan, if any.
620    pub fn pending_followup_draft(&self) -> Option<&str> {
621        self.pending_followup_draft.as_deref()
622    }
623
624    /// Insert the launch plan's deferred draft exactly once without submitting.
625    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    /// Submit the launch plan's deferred prompt exactly once.
647    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    /// Send a prompt char-by-char (required for Ink-based TUI tools like Claude Code).
697    ///
698    /// Sends each character with a small delay to avoid overwhelming the TUI's raw-mode
699    /// input processing. Ends with a carriage return.
700    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    /// Resize the PTY.
716    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    /// Deliver a platform-neutral PTY control or request process termination.
732    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    /// Session ID assigned at spawn time.
766    pub fn session_id(&self) -> &str {
767        &self.session_id
768    }
769
770    /// Current event-stream generation.
771    pub fn generation(&self) -> u64 {
772        self.events.generation()
773    }
774
775    /// Revision pinned into event envelopes, cursors, and snapshots.
776    pub fn provider_revision(&self) -> &str {
777        self.events.provider_revision()
778    }
779
780    /// Cursor for replay from the beginning of this generation.
781    pub fn beginning_cursor(&self) -> PtyReplayCursor {
782        PtyReplayCursor::beginning(self.provider_revision(), self.generation())
783    }
784
785    /// Cursor for the earliest event still retained by the bounded journal.
786    /// Long-lived readiness probes use this to build fresh positive evidence
787    /// from the available tail without treating normal eviction as corruption.
788    pub fn retained_cursor(&self) -> Result<PtyReplayCursor, PtyAttachError> {
789        self.events.retained_cursor()
790    }
791
792    /// Atomically replay the currently retained journal tail and subscribe
793    /// after its exact boundary.
794    pub fn attach_retained_events(&self) -> Result<PtyAttachment, PtyAttachError> {
795        self.events.attach_retained()
796    }
797
798    /// Stable agent identity associated with the session.
799    pub fn agent_id(&self) -> &AgentId {
800        &self.agent_id
801    }
802
803    /// Root process ID observed from the owned PTY child handle.
804    pub fn root_pid(&self) -> Option<u32> {
805        self.root_pid
806    }
807
808    /// Whether the output reader has reached a terminal state.
809    pub fn reader_finished(&self) -> bool {
810        self.reader_task
811            .as_ref()
812            .is_none_or(JoinHandle::is_finished)
813    }
814
815    /// Snapshot and terminate agent descendants before killing the root handle.
816    /// The returned report makes any root-only degradation explicit.
817    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    /// Kill the process tree. The reader remains alive long enough to drain
839    /// accepted PTY output and publish the ordered exit event.
840    pub async fn kill(&self) -> Result<(), AgentError> {
841        self.terminate_tree().await.map(|_| ())
842    }
843
844    /// Consume the session, terminate owned processes, and join the ordered
845    /// reader/exit path within a bounded deadline.
846    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                // The kill itself may well have landed (`termination`,
862                // above) -- this is the reader thread failing to join and
863                // observe that within budget, a separate and, until now,
864                // silent way for a PTY shutdown to get stuck.
865                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
916// ---------------------------------------------------------------------------
917// Reader loop (runs on blocking thread via spawn_blocking)
918// ---------------------------------------------------------------------------
919
920fn 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        // Broadcast raw PTY bytes (for vt100 screen emulation by consumers).
1010        // Keep as Vec<u8> — never round-trip through String::from_utf8_lossy
1011        // because that destroys multi-byte UTF-8 sequences split across reads.
1012        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        // For text analysis (rate limit detection, classification) we need a
1020        // String. Lossy conversion is acceptable here because this branch only
1021        // feeds the heuristic ANSI/text parsers, not the vt100 grid.
1022        let raw_str = String::from_utf8_lossy(&raw).to_string();
1023
1024        // Strip ANSI for text analysis
1025        let cleaned = vte_parser.parse(&raw_str);
1026
1027        // Rate limit detection
1028        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        // Classification pipeline
1036        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        // Startup sequence handling
1059        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
1086/// Generate a simple UUID-like session ID.
1087fn 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}