diff --git a/src/pane.rs b/src/pane.rs index c53d99ef..baa5790c 100644 --- a/src/pane.rs +++ b/src/pane.rs @@ -118,7 +118,7 @@ const PROCESS_ACQUISITION_FAST_RECHECK: std::time::Duration = std::time::Duratio const PROCESS_ACQUISITION_SLOW_RECHECK: std::time::Duration = std::time::Duration::from_secs(2); const PROCESS_ACQUISITION_IDLE_RESET: std::time::Duration = std::time::Duration::from_secs(2); const STABLE_VISIBLE_SIGNAL_REFRESH: std::time::Duration = std::time::Duration::from_millis(800); -const AGENT_PTY_ACTIVITY_WINDOW: std::time::Duration = std::time::Duration::from_secs(1); +const AGENT_PTY_ACTIVITY_WINDOW: std::time::Duration = std::time::Duration::from_millis(1800); const AGENT_INPUT_TAINT_WINDOW: std::time::Duration = std::time::Duration::from_millis(1200); const AGENT_STARTUP_GRACE_WINDOW: std::time::Duration = std::time::Duration::from_secs(3); const AGENT_PENDING_IDLE_RECHECK: std::time::Duration = std::time::Duration::from_millis(100); @@ -366,6 +366,7 @@ impl PendingIdleConfirmation { next: DetectionPublishState, agent_changed: bool, process_exited: bool, + pty_signal: Option, now: std::time::Instant, ) -> bool { let is_working_to_plain_idle = previous.state == AgentState::Working @@ -379,6 +380,11 @@ impl PendingIdleConfirmation { return false; } + if pty_signal.is_some_and(|signal| !signal.active && !signal.tainted) { + self.clear(); + return false; + } + let Some(started_at) = self.started_at else { self.started_at = Some(now); self.confirmations = 0; @@ -403,7 +409,7 @@ impl PendingIdleConfirmation { #[derive(Debug, Default)] struct PendingWorkingConfirmation { started_at: Option, - last_observed_render_seq: u64, + last_observed_output_seq: u64, } impl PendingWorkingConfirmation { @@ -413,7 +419,7 @@ impl PendingWorkingConfirmation { fn clear(&mut self) { self.started_at = None; - self.last_observed_render_seq = 0; + self.last_observed_output_seq = 0; } fn recheck_delay(&self, now: std::time::Instant) -> std::time::Duration { @@ -458,12 +464,12 @@ impl PendingWorkingConfirmation { let Some(started_at) = self.started_at else { self.started_at = Some(now); - self.last_observed_render_seq = pty_signal.render_seq; + self.last_observed_output_seq = pty_signal.output_seq; return true; }; - if pty_signal.fresh_render && pty_signal.render_seq != self.last_observed_render_seq { - self.last_observed_render_seq = pty_signal.render_seq; + if pty_signal.fresh_output && pty_signal.output_seq != self.last_observed_output_seq { + self.last_observed_output_seq = pty_signal.output_seq; if now.duration_since(started_at) >= AGENT_PENDING_WORKING_CONFIRM_DELAY { self.clear(); return false; @@ -562,7 +568,7 @@ fn detection_update_for_publish( #[derive(Debug, Default)] struct PtyCausalityTracker { - last_pty_render_seq: u64, + last_pty_output_seq: u64, last_input_seq: u64, input_tainted_until: Option, last_agent_pty_at: Option, @@ -572,19 +578,19 @@ struct PtyCausalityTracker { struct PtyActivitySignal { active: bool, tainted: bool, - fresh_render: bool, - render_seq: u64, + fresh_output: bool, + output_seq: u64, } -fn baseline_pty_causality(tracker: &mut PtyCausalityTracker, pty_render_seq: u64, input_seq: u64) { - tracker.last_pty_render_seq = pty_render_seq; +fn baseline_pty_causality(tracker: &mut PtyCausalityTracker, pty_output_seq: u64, input_seq: u64) { + tracker.last_pty_output_seq = pty_output_seq; tracker.last_input_seq = input_seq; tracker.input_tainted_until = None; tracker.last_agent_pty_at = None; } fn agent_caused_pty_activity_active( - pty_render_seq: u64, + pty_output_seq: u64, input_seq: u64, tracker: &mut PtyCausalityTracker, now: std::time::Instant, @@ -597,12 +603,12 @@ fn agent_caused_pty_activity_active( let tainted = tracker.input_tainted_until.is_some_and(|until| now < until); - let mut fresh_render = false; - if pty_render_seq != tracker.last_pty_render_seq { - tracker.last_pty_render_seq = pty_render_seq; + let mut fresh_output = false; + if pty_output_seq != tracker.last_pty_output_seq { + tracker.last_pty_output_seq = pty_output_seq; if !tainted { tracker.last_agent_pty_at = Some(now); - fresh_render = true; + fresh_output = true; } } @@ -610,8 +616,8 @@ fn agent_caused_pty_activity_active( return PtyActivitySignal { active: false, tainted: true, - fresh_render: false, - render_seq: tracker.last_pty_render_seq, + fresh_output: false, + output_seq: tracker.last_pty_output_seq, }; } @@ -621,8 +627,14 @@ fn agent_caused_pty_activity_active( PtyActivitySignal { active, tainted: false, - fresh_render, - render_seq: tracker.last_pty_render_seq, + fresh_output, + output_seq: tracker.last_pty_output_seq, + } +} + +fn observe_pty_output_activity(bytes: &[u8], pty_output_seq: &AtomicU64) { + if !bytes.is_empty() { + pty_output_seq.fetch_add(1, Ordering::Relaxed); } } @@ -631,7 +643,7 @@ fn spawn_basic_detection_task( pane_id: PaneId, child_pid: Arc, terminal: Arc, - pty_render_seq: Arc, + pty_output_seq: Arc, input_write_seq: Arc, state_events: mpsc::Sender, ) -> ( @@ -786,7 +798,7 @@ fn spawn_basic_detection_task( agent_startup_grace_until = Some(now + AGENT_STARTUP_GRACE_WINDOW); baseline_pty_causality( &mut pty_causality, - pty_render_seq.load(Ordering::Relaxed), + pty_output_seq.load(Ordering::Relaxed), input_write_seq.load(Ordering::Relaxed), ); state = AgentState::Idle; @@ -828,7 +840,7 @@ fn spawn_basic_detection_task( } baseline_pty_causality( &mut pty_causality, - pty_render_seq.load(Ordering::Relaxed), + pty_output_seq.load(Ordering::Relaxed), input_write_seq.load(Ordering::Relaxed), ); agent_startup_grace_until = None; @@ -865,7 +877,7 @@ fn spawn_basic_detection_task( }; let pty_activity = if agent.is_some() { Some(agent_caused_pty_activity_active( - pty_render_seq.load(Ordering::Relaxed), + pty_output_seq.load(Ordering::Relaxed), input_write_seq.load(Ordering::Relaxed), &mut pty_causality, now, @@ -967,6 +979,7 @@ fn spawn_basic_detection_task( next_publish, agent_changed, process_exited, + pty_activity, now, ) { pending_working.clear(); @@ -1817,25 +1830,25 @@ impl PaneRuntime { let reported_cwd = Arc::new(Mutex::new(None)); let kitty_keyboard_flags = Arc::new(AtomicU16::new(keyboard_protocol_flags)); let input_write_seq = Arc::new(AtomicU64::new(0)); - let pty_render_seq = Arc::new(AtomicU64::new(0)); + let pty_output_seq = Arc::new(AtomicU64::new(0)); let io = { let terminal = terminal.clone(); let response_writer = response_tx.clone(); let render_notify = render_notify.clone(); let render_dirty = render_dirty.clone(); - let pty_render_seq = pty_render_seq.clone(); + let pty_output_seq = pty_output_seq.clone(); let child_pid = child_pid.clone(); let read_events = events.clone(); let reported_cwd = reported_cwd.clone(); let rt = tokio::runtime::Handle::current(); let delay_rt = rt.clone(); let on_read = Box::new(move |bytes: &[u8]| { + observe_pty_output_activity(bytes, &pty_output_seq); let shell_pid = child_pid.load(Ordering::Acquire); let result = terminal.process_pty_bytes(pane_id, shell_pid, bytes, &response_writer); if result.request_render { - pty_render_seq.fetch_add(1, Ordering::Relaxed); if !render_dirty.swap(true, Ordering::AcqRel) { render_notify.notify_one(); } @@ -1884,7 +1897,7 @@ impl PaneRuntime { pane_id, child_pid.clone(), terminal.clone(), - pty_render_seq, + pty_output_seq, input_write_seq.clone(), events, ); @@ -1948,7 +1961,7 @@ impl PaneRuntime { let reported_cwd = Arc::new(Mutex::new(None)); let child_wait_completed = Arc::new(AtomicBool::new(false)); let input_write_seq = Arc::new(AtomicU64::new(0)); - let pty_render_seq = Arc::new(AtomicU64::new(0)); + let pty_output_seq = Arc::new(AtomicU64::new(0)); { let child_pid = child_pid.clone(); let child_wait_completed = child_wait_completed.clone(); @@ -1980,17 +1993,17 @@ impl PaneRuntime { let response_writer = response_tx.clone(); let render_notify = render_notify.clone(); let render_dirty = render_dirty.clone(); - let pty_render_seq = pty_render_seq.clone(); + let pty_output_seq = pty_output_seq.clone(); let child_pid = child_pid.clone(); let events = events.clone(); let reported_cwd = reported_cwd.clone(); let rt = tokio::runtime::Handle::current(); let on_read = Box::new(move |bytes: &[u8]| { + observe_pty_output_activity(bytes, &pty_output_seq); let shell_pid = child_pid.load(Ordering::Acquire); let result = terminal.process_pty_bytes(pane_id, shell_pid, bytes, &response_writer); if result.request_render { - pty_render_seq.fetch_add(1, Ordering::Relaxed); if !render_dirty.swap(true, Ordering::AcqRel) { render_notify.notify_one(); } @@ -2045,7 +2058,7 @@ impl PaneRuntime { let child_pid = child_pid.clone(); let terminal = terminal.clone(); let state_events = events.clone(); - let pty_render_seq = pty_render_seq.clone(); + let pty_output_seq = pty_output_seq.clone(); let input_write_seq_for_task = input_write_seq.clone(); let render_notify = render_notify.clone(); let render_dirty = render_dirty.clone(); @@ -2215,7 +2228,7 @@ impl PaneRuntime { Some(now + AGENT_STARTUP_GRACE_WINDOW); baseline_pty_causality( &mut pty_causality, - pty_render_seq.load(Ordering::Relaxed), + pty_output_seq.load(Ordering::Relaxed), input_write_seq_for_task.load(Ordering::Relaxed), ); state = AgentState::Idle; @@ -2286,7 +2299,7 @@ impl PaneRuntime { } baseline_pty_causality( &mut pty_causality, - pty_render_seq.load(Ordering::Relaxed), + pty_output_seq.load(Ordering::Relaxed), input_write_seq_for_task.load(Ordering::Relaxed), ); agent_startup_grace_until = None; @@ -2323,7 +2336,7 @@ impl PaneRuntime { }; let pty_activity = if agent.is_some() { Some(agent_caused_pty_activity_active( - pty_render_seq.load(Ordering::Relaxed), + pty_output_seq.load(Ordering::Relaxed), input_write_seq_for_task.load(Ordering::Relaxed), &mut pty_causality, now, @@ -2432,6 +2445,7 @@ impl PaneRuntime { next_publish, agent_changed, process_exited, + pty_activity, now, ) { pending_working.clear(); @@ -3416,12 +3430,13 @@ mod tests { }; let mut pending = PendingIdleConfirmation::default(); - assert!(pending.should_hold_working_to_idle(previous, next, false, false, now)); + assert!(pending.should_hold_working_to_idle(previous, next, false, false, None, now)); assert!(pending.should_hold_working_to_idle( previous, next, false, false, + None, now + AGENT_PENDING_IDLE_RECHECK )); assert!(pending.should_hold_working_to_idle( @@ -3429,6 +3444,7 @@ mod tests { next, false, false, + None, now + AGENT_PENDING_IDLE_RECHECK * 2 )); assert!(!pending.should_hold_working_to_idle( @@ -3436,6 +3452,7 @@ mod tests { next, false, false, + None, now + AGENT_PENDING_IDLE_RECHECK * 3 )); } @@ -3455,12 +3472,13 @@ mod tests { }; let mut pending = PendingIdleConfirmation::default(); - assert!(pending.should_hold_working_to_idle(previous, next, false, false, now)); + assert!(pending.should_hold_working_to_idle(previous, next, false, false, None, now)); assert!(!pending.should_hold_working_to_idle( previous, next, false, false, + None, now + AGENT_PENDING_IDLE_CAP )); } @@ -3490,22 +3508,48 @@ mod tests { }; let mut pending = PendingIdleConfirmation::default(); - assert!(pending.should_hold_working_to_idle(previous, idle, false, false, now)); + assert!(pending.should_hold_working_to_idle(previous, idle, false, false, None, now)); assert!(pending.active()); - assert!(!pending.should_hold_working_to_idle(previous, working, false, false, now)); + assert!(!pending.should_hold_working_to_idle(previous, working, false, false, None, now)); assert!(!pending.active()); - assert!(pending.should_hold_working_to_idle(previous, idle, false, false, now)); - assert!(!pending.should_hold_working_to_idle(previous, blocked, false, false, now)); + assert!(pending.should_hold_working_to_idle(previous, idle, false, false, None, now)); + assert!(!pending.should_hold_working_to_idle(previous, blocked, false, false, None, now)); assert!(!pending.active()); } - fn pty_activity(active: bool, fresh_render: bool, render_seq: u64) -> PtyActivitySignal { + #[test] + fn pending_idle_does_not_extend_pty_quiet_lease() { + let now = std::time::Instant::now(); + let previous = DetectionPublishState { + state: AgentState::Working, + visible_blocker: false, + visible_working: false, + }; + let idle = DetectionPublishState { + state: AgentState::Idle, + visible_blocker: false, + visible_working: false, + }; + let mut pending = PendingIdleConfirmation::default(); + + assert!(!pending.should_hold_working_to_idle( + previous, + idle, + false, + false, + Some(pty_activity(false, false, 10)), + now + )); + assert!(!pending.active()); + } + + fn pty_activity(active: bool, fresh_output: bool, output_seq: u64) -> PtyActivitySignal { PtyActivitySignal { active, tainted: false, - fresh_render, - render_seq, + fresh_output, + output_seq, } } @@ -3545,7 +3589,7 @@ mod tests { } #[test] - fn pending_working_confirms_after_delayed_render_observation() { + fn pending_working_confirms_after_delayed_output_observation() { let now = std::time::Instant::now(); let idle = DetectionPublishState { state: AgentState::Idle, @@ -3715,15 +3759,15 @@ mod tests { } #[test] - fn pty_activity_distinguishes_fresh_render_from_active_hold() { + fn pty_activity_distinguishes_fresh_output_from_active_hold() { let now = std::time::Instant::now(); let mut tracker = PtyCausalityTracker::default(); baseline_pty_causality(&mut tracker, 1, 1); let fresh = agent_caused_pty_activity_active(2, 1, &mut tracker, now); assert!(fresh.active); - assert!(fresh.fresh_render); - assert_eq!(fresh.render_seq, 2); + assert!(fresh.fresh_output); + assert_eq!(fresh.output_seq, 2); let held = agent_caused_pty_activity_active( 2, @@ -3732,8 +3776,67 @@ mod tests { now + AGENT_PENDING_WORKING_FAST_RECHECK, ); assert!(held.active); - assert!(!held.fresh_render); - assert_eq!(held.render_seq, 2); + assert!(!held.fresh_output); + assert_eq!(held.output_seq, 2); + } + + #[test] + fn pty_activity_lease_survives_sparse_heartbeat_jitter() { + let now = std::time::Instant::now(); + let mut tracker = PtyCausalityTracker::default(); + baseline_pty_causality(&mut tracker, 1, 1); + + let first = agent_caused_pty_activity_active(2, 1, &mut tracker, now); + assert!(first.active); + assert!(first.fresh_output); + + let held_before_next_tick = agent_caused_pty_activity_active( + 2, + 1, + &mut tracker, + now + std::time::Duration::from_millis(1700), + ); + assert!(held_before_next_tick.active); + assert!(!held_before_next_tick.fresh_output); + + let next_tick = agent_caused_pty_activity_active( + 3, + 1, + &mut tracker, + now + std::time::Duration::from_millis(1700), + ); + assert!(next_tick.active); + assert!(next_tick.fresh_output); + + let held_after_next_tick = agent_caused_pty_activity_active( + 3, + 1, + &mut tracker, + now + std::time::Duration::from_millis(3400), + ); + assert!(held_after_next_tick.active); + + let expired = agent_caused_pty_activity_active( + 3, + 1, + &mut tracker, + now + std::time::Duration::from_millis(3501), + ); + assert!(!expired.active); + } + + #[test] + fn pty_output_activity_tracks_raw_nonempty_reads() { + let seq = AtomicU64::new(0); + + observe_pty_output_activity(b"", &seq); + assert_eq!(seq.load(Ordering::Relaxed), 0); + + observe_pty_output_activity(b"\x1b[?2026h", &seq); + assert_eq!(seq.load(Ordering::Relaxed), 1); + + observe_pty_output_activity(b"body bytes", &seq); + assert_eq!(seq.load(Ordering::Relaxed), 2); } #[test] @@ -4178,7 +4281,7 @@ mod tests { } #[test] - fn input_taint_discards_pty_activity_until_fresh_post_taint_render() { + fn input_taint_discards_pty_activity_until_fresh_post_taint_output() { let now = std::time::Instant::now(); let mut tracker = PtyCausalityTracker::default(); baseline_pty_causality(&mut tracker, 1, 1); @@ -4196,14 +4299,14 @@ mod tests { assert!(!after_taint.active); assert!(!after_taint.tainted); - let fresh_render = agent_caused_pty_activity_active( + let fresh_output = agent_caused_pty_activity_active( 3, 2, &mut tracker, now + AGENT_INPUT_TAINT_WINDOW + std::time::Duration::from_millis(2), ); - assert!(fresh_render.active); - assert!(!fresh_render.tainted); + assert!(fresh_output.active); + assert!(!fresh_output.tainted); } #[test]