Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions codex-rs/tui/src/app.rs
Original file line number Diff line number Diff line change
Expand Up @@ -229,6 +229,7 @@ mod session_lifecycle;
mod side;
mod startup;
mod startup_prompts;
mod thread_event_buffer;
mod thread_events;
mod thread_goal_actions;
mod thread_routing;
Expand Down
12 changes: 1 addition & 11 deletions codex-rs/tui/src/app/background_requests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -653,17 +653,7 @@ impl App {

let should_send = {
let mut guard = store.lock().await;
guard
.buffer
.push_back(ThreadBufferedEvent::FeedbackSubmission(event.clone()));
if guard.buffer.len() > guard.capacity
&& let Some(removed) = guard.buffer.pop_front()
&& let ThreadBufferedEvent::Request(request) = &removed
{
guard
.pending_interactive_replay
.note_evicted_server_request(request.as_ref());
}
guard.push_buffered_event(ThreadBufferedEvent::FeedbackSubmission(event.clone()));
guard.active
};

Expand Down
183 changes: 183 additions & 0 deletions codex-rs/tui/src/app/tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -743,6 +743,81 @@ async fn active_history_batch_is_delivered_without_replay_buffering() -> Result<
Ok(())
}

#[tokio::test]
async fn replay_thread_snapshot_renders_only_retained_agent_message_deltas() {
let (mut app, mut app_event_rx, _op_rx) = make_test_app_with_channels().await;
let thread_id = ThreadId::new();
let session = test_thread_session(thread_id, test_path_buf("/tmp/project"));
app.thread_event_channels.insert(
thread_id,
ThreadEventChannel::new_with_session(
THREAD_EVENT_CHANNEL_CAPACITY,
session.clone(),
Vec::new(),
),
);
app.activate_thread_channel(thread_id).await;
app.chat_widget.handle_thread_session(session);

{
let channel = app
.thread_event_channels
.get(&thread_id)
.expect("thread channel should exist");
let mut store = channel.store.lock().await;
for chunk in 0..65 {
let marker = format!("chunk {chunk:02}");
let text = format!("{marker}\n{}\n", "x".repeat(4096 - marker.len() - 2));
store.push_notification(agent_message_delta_notification(
thread_id,
"turn-budget",
"item-budget",
&text,
));
}
}

app.store_active_thread_receiver().await;
let (_receiver, snapshot) = app
.activate_thread_for_replay(thread_id)
.await
.expect("detached thread should reactivate for replay");
app.replay_thread_snapshot(snapshot, /*resume_restored_queue*/ false);

let mut rendered = Vec::new();
loop {
app.chat_widget.on_commit_tick();
let mut inserted = false;
while let Ok(event) = app_event_rx.try_recv() {
if let AppEvent::InsertHistoryCell(cell) = event {
rendered.push(lines_to_single_string(&cell.display_lines(/*width*/ 80)));
inserted = true;
}
}
if !inserted {
break;
}
}
if let Some(lines) = app.chat_widget.active_cell_transcript_lines(/*width*/ 80) {
rendered.push(lines_to_single_string(&lines));
}
let markers = rendered
.iter()
.flat_map(|text| text.lines())
.filter_map(|line| line.split_once("chunk "))
.map(|(_, marker)| format!("chunk {}", marker.trim()))
.collect::<Vec<_>>();

assert_eq!(markers.len(), 64);
assert_snapshot!(
format!("{}\n{}", markers.first().expect("first marker"), markers.last().expect("last marker")),
@r"
chunk 01
chunk 64
"
);
}

#[tokio::test]
async fn replay_thread_snapshot_restores_draft_and_queued_input() {
let mut app = make_test_app().await;
Expand Down Expand Up @@ -4111,6 +4186,114 @@ async fn side_parent_status_prioritizes_input_over_approval() -> Result<()> {
None
);

app.enqueue_thread_request(
parent_thread_id,
request_user_input_request(parent_thread_id, "turn-eviction", "input-eviction"),
)
.await?;
let chunk = "x".repeat(4 * 1024);
app.enqueue_thread_notification(
parent_thread_id,
agent_message_delta_notification(
parent_thread_id,
"turn-eviction",
"item-eviction",
&chunk,
),
)
.await?;
app.enqueue_thread_request(
parent_thread_id,
exec_approval_request(
parent_thread_id,
"turn-eviction",
"approval-eviction",
/*approval_id*/ None,
),
)
.await?;

let side_footer = |app: &App| {
render_bottom_popup(&app.chat_widget, /*width*/ 120)
.lines()
.find_map(|line| {
line.find("Side from main thread")
.map(|start| line[start..].trim().to_string())
})
.expect("side conversation footer should be rendered")
};
assert_eq!(
app.side_threads
.get(&side_thread_id)
.and_then(|state| state.parent_status),
Some(SideParentStatus::NeedsInput)
);
let input_footer = side_footer(&app);

for _ in 0..64 {
app.enqueue_thread_notification(
parent_thread_id,
agent_message_delta_notification(
parent_thread_id,
"turn-eviction",
"item-eviction",
&chunk,
),
)
.await?;
}
assert_eq!(
app.side_threads
.get(&side_thread_id)
.and_then(|state| state.parent_status),
Some(SideParentStatus::NeedsApproval)
);
let approval_footer = side_footer(&app);
assert!(
app.thread_event_channels
.get(&parent_thread_id)
.expect("parent thread channel should exist")
.store
.lock()
.await
.has_pending_thread_approvals()
);

app.enqueue_thread_notification(
parent_thread_id,
agent_message_delta_notification(
parent_thread_id,
"turn-eviction",
"item-eviction",
&chunk,
),
)
.await?;
assert_eq!(
app.side_threads
.get(&side_thread_id)
.and_then(|state| state.parent_status),
None
);
let cleared_footer = side_footer(&app);
assert!(
!app.thread_event_channels
.get(&parent_thread_id)
.expect("parent thread channel should exist")
.store
.lock()
.await
.has_pending_thread_approvals()
);
assert_snapshot!(
format!("{input_footer}\n{approval_footer}\n{cleared_footer}"),
@r"
Side from main thread · main needs input · ctrl + / to switch · ctrl + c to close
Side from main thread · main needs approval · ctrl + / to switch · ctrl + c to close
Side from main thread · ctrl + / to switch · ctrl + c to close
"
);

Ok(())
}

Expand Down
81 changes: 81 additions & 0 deletions codex-rs/tui/src/app/thread_event_buffer.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,81 @@
//! Bounded replay-buffer policy for per-thread app-server events.

use super::ServerNotification;
use super::ThreadBufferedEvent;
use super::ThreadEventStore;
use std::borrow::Cow;

// Keep merged text finite so continued streaming still reaches bounded replay eviction.
const MAX_COALESCED_AGENT_MESSAGE_DELTA_BYTES: usize = 4 * 1024;
const MAX_BUFFERED_AGENT_MESSAGE_DELTA_BYTES: usize = 256 * 1024;

impl ThreadEventStore {
pub(super) fn push_replay_notification(&mut self, notification: Cow<'_, ServerNotification>) {
if let ServerNotification::AgentMessageDelta(delta) = notification.as_ref()
&& delta.delta.len() > MAX_BUFFERED_AGENT_MESSAGE_DELTA_BYTES
{
return;
}

if let ServerNotification::AgentMessageDelta(delta) = notification.as_ref()
&& let Some(ThreadBufferedEvent::Notification(previous)) = self.buffer.back_mut()
&& let ServerNotification::AgentMessageDelta(previous) = previous.as_mut()
&& previous.thread_id == delta.thread_id
&& previous.turn_id == delta.turn_id
&& previous.item_id == delta.item_id
&& previous.delta.len().saturating_add(delta.delta.len())
<= MAX_COALESCED_AGENT_MESSAGE_DELTA_BYTES
{
previous.delta.push_str(&delta.delta);
self.buffered_agent_message_delta_bytes = self
.buffered_agent_message_delta_bytes
.saturating_add(delta.delta.len());
self.evict_overflowing_events();
return;
}

self.push_buffered_event(ThreadBufferedEvent::Notification(Box::new(
notification.into_owned(),
)));
}

pub(super) fn push_buffered_event(&mut self, event: ThreadBufferedEvent) {
if let ThreadBufferedEvent::Notification(notification) = &event
&& let ServerNotification::AgentMessageDelta(delta) = notification.as_ref()
{
self.buffered_agent_message_delta_bytes = self
.buffered_agent_message_delta_bytes
.saturating_add(delta.delta.len());
}
self.buffer.push_back(event);
self.evict_overflowing_events();
}

fn evict_overflowing_events(&mut self) {
while self.buffer.len() > self.capacity
|| self.buffered_agent_message_delta_bytes > MAX_BUFFERED_AGENT_MESSAGE_DELTA_BYTES
{
let Some(removed) = self.buffer.pop_front() else {
break;
};
match removed {
ThreadBufferedEvent::Notification(notification) => {
if let ServerNotification::AgentMessageDelta(delta) = notification.as_ref() {
self.buffered_agent_message_delta_bytes = self
.buffered_agent_message_delta_bytes
.saturating_sub(delta.delta.len());
}
}
ThreadBufferedEvent::Request(request) => self
.pending_interactive_replay
.note_evicted_server_request(request.as_ref()),
ThreadBufferedEvent::HistoryEntryResponse(_)
| ThreadBufferedEvent::FeedbackSubmission(_) => {}
}
}
}
}

#[cfg(test)]
#[path = "thread_event_buffer_tests.rs"]
mod tests;
Loading
Loading