fix: lease working status on pty output

This commit is contained in:
Ogulcan Celik
2026-06-09 04:26:19 +03:00
parent 615310e15c
commit d5ccf19a9d
+157 -54
View File
@@ -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<PtyActivitySignal>,
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<std::time::Instant>,
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<std::time::Instant>,
last_agent_pty_at: Option<std::time::Instant>,
@@ -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<AtomicU32>,
terminal: Arc<PaneTerminal>,
pty_render_seq: Arc<AtomicU64>,
pty_output_seq: Arc<AtomicU64>,
input_write_seq: Arc<AtomicU64>,
state_events: mpsc::Sender<AppEvent>,
) -> (
@@ -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]