use super::*;
use app_test_support::create_fake_parented_rollout_with_source;
use app_test_support::create_fake_rollout;
use app_test_support::rollout_path;
use lemurclaw_core::app_server_protocol::ClientNotification;
use lemurclaw_core::app_server_protocol::ClientRequest;
use lemurclaw_core::app_server_protocol::JSONRPCError;
use lemurclaw_core::app_server_protocol::JSONRPCMessage;
use lemurclaw_core::app_server_protocol::JSONRPCResponse;
use lemurclaw_core::protocol::AgentPath;
use futures::SinkExt;
use futures::StreamExt;
use pretty_assertions::assert_eq;
use std::sync::Mutex;
use tokio::net::TcpListener;
use tokio::task::JoinHandle;
use tokio_tungstenite::accept_async;
use tokio_tungstenite::tungstenite::Message;
fn take_backfill_counts(requests: &Arc<Mutex<Vec<String>>>) -> (usize, usize) {
let requests = std::mem::take(&mut *requests.lock().expect("request recorder lock"));
(
requests
.iter()
.filter(|method| *method == "thread/loaded/list")
.count(),
requests
.iter()
.filter(|method| *method == "thread/read")
.count(),
)
}
async fn start_recording_app_server(
config: &Config,
) -> Result<(
AppServerSession,
Arc<Mutex<Vec<String>>>,
JoinHandle<Result<()>>,
)> {
let state_db =
crate::tui_internal::init_state_db_for_app_server_target(config, &crate::tui_internal::AppServerTarget::Embedded)
.await?;
let embedded = crate::tui_internal::start_embedded_app_server(
lemurclaw_server::arg0::Arg0DispatchPaths::default(),
config.clone(),
Vec::new(),
lemurclaw_core::config::LoaderOverrides::default(),
false,
lemurclaw_core::config::CloudConfigBundleLoader::default(),
lemurclaw_core::feedback::CodexFeedback::new(),
None,
state_db,
Arc::new(lemurclaw_core::exec_server::EnvironmentManager::default_for_tests()),
)
.await?;
let codex_home = config.codex_home.display().to_string();
let requests = Arc::new(Mutex::new(Vec::new()));
let request_sink = Arc::clone(&requests);
let listener = TcpListener::bind("127.0.0.1:0").await?;
let websocket_url = format!("ws://{}", listener.local_addr()?);
let proxy = tokio::spawn(async move {
let (stream, _) = listener.accept().await?;
let mut websocket = accept_async(stream).await?;
while let Some(frame) = websocket.next().await {
let Message::Text(text) = frame? else {
continue;
};
let message = serde_json::from_str::<JSONRPCMessage>(&text)?;
match message {
JSONRPCMessage::Request(request) if request.method == "initialize" => {
websocket
.send(Message::Text(
serde_json::to_string(&JSONRPCMessage::Response(JSONRPCResponse {
id: request.id,
result: serde_json::json!({
"userAgent": "codex-tui-test",
"codexHome": codex_home,
}),
}))?
.into(),
))
.await?;
}
JSONRPCMessage::Request(request) => {
request_sink
.lock()
.expect("request recorder lock")
.push(request.method.clone());
let request_id = request.id.clone();
let request =
serde_json::from_value::<ClientRequest>(serde_json::to_value(request)?)?;
let response = match embedded.request(request).await? {
Ok(result) => JSONRPCMessage::Response(JSONRPCResponse {
id: request_id,
result,
}),
Err(error) => JSONRPCMessage::Error(JSONRPCError {
id: request_id,
error,
}),
};
websocket
.send(Message::Text(serde_json::to_string(&response)?.into()))
.await?;
}
JSONRPCMessage::Notification(notification)
if notification.method == "initialized" => {}
JSONRPCMessage::Notification(notification) => {
embedded
.notify(serde_json::from_value::<ClientNotification>(
serde_json::to_value(notification)?,
)?)
.await?;
}
JSONRPCMessage::Response(_) | JSONRPCMessage::Error(_) => {}
}
}
embedded.shutdown().await?;
Ok(())
});
let app_server = crate::tui_internal::connect_remote_app_server(crate::tui_internal::RemoteAppServerEndpoint::WebSocket {
websocket_url,
auth_token: None,
})
.await?;
Ok((
AppServerSession::new(
app_server,
crate::tui_internal::app_server_session::ThreadParamsMode::Embedded,
),
requests,
proxy,
))
}
#[test]
fn session_lifecycle_avoids_redundant_subagent_metadata_reads() -> Result<()> {
const TEST_STACK_SIZE_BYTES: usize = 8 * 1024 * 1024;
std::thread::Builder::new()
.name("tui-session-lifecycle-requests".to_string())
.stack_size(TEST_STACK_SIZE_BYTES)
.spawn(|| {
let runtime = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()?;
runtime.block_on(async {
let mut app = make_test_app().await;
let codex_home = tempdir()?;
app.config.codex_home = codex_home.path().to_path_buf().abs();
app.config.sqlite_home = codex_home.path().to_path_buf();
let root_timestamp = "2026-01-01T00-00-00";
let root_thread_id = ThreadId::from_string(
&create_fake_rollout(
codex_home.path(),
root_timestamp,
"2026-01-01T00:00:00Z",
"Saved user message",
Some(app.config.model_provider_id.as_str()),
None,
)
.expect("create root rollout"),
)?;
let child_thread_id = ThreadId::from_string(
&create_fake_parented_rollout_with_source(
codex_home.path(),
"2026-01-01T00-00-01",
"2026-01-01T00:00:01Z",
"Saved child message",
Some(app.config.model_provider_id.as_str()),
None,
RolloutSessionSource::SubAgent(SubAgentSource::ThreadSpawn {
parent_thread_id: root_thread_id,
depth: 1,
agent_path: Some(
AgentPath::try_from("/root/worker").expect("valid agent path"),
),
agent_nickname: Some("worker".to_string()),
agent_role: Some("worker".to_string()),
}),
root_thread_id.into(),
root_thread_id,
)
.expect("create child rollout"),
)?;
let root_rollout_path = rollout_path(
codex_home.path(),
root_timestamp,
&root_thread_id.to_string(),
);
let (mut app_server, requests, proxy) =
start_recording_app_server(&app.config).await?;
let root = app_server
.resume_thread(
app.config.clone(),
root_thread_id,
app.resume_model_settings(),
)
.await?;
app.enqueue_primary_thread_session(root.session, root.turns)
.await?;
app_server
.resume_thread(
app.config.clone(),
child_thread_id,
app.resume_model_settings(),
)
.await?;
let mut tui = crate::tui_internal::tui::test_support::make_test_tui()?;
take_backfill_counts(&requests);
let control = Box::pin(app.handle_event(
&mut tui,
&mut app_server,
AppEvent::ForkCurrentSession,
))
.await?;
assert!(matches!(control, AppRunControl::Continue));
assert_ne!(app.chat_widget.thread_id(), Some(root_thread_id));
assert!(matches!(take_backfill_counts(&requests), (0, 0) | (0, 1)));
app.start_fresh_session_with_summary_hint(
&mut tui,
&mut app_server,
None,
None,
)
.await;
assert_ne!(app.chat_widget.thread_id(), Some(root_thread_id));
assert_eq!(take_backfill_counts(&requests), (0, 0));
let loaded_threads = app_server
.thread_loaded_list(ThreadLoadedListParams {
cursor: None,
limit: None,
})
.await?
.data;
let expected_reads = loaded_threads
.iter()
.filter(|thread_id| *thread_id != &root_thread_id.to_string())
.count();
assert!(loaded_threads.contains(&child_thread_id.to_string()));
take_backfill_counts(&requests);
app.harness_overrides.cwd = Some(app.config.cwd.to_path_buf());
let control = app
.resume_target_session(
&mut tui,
&mut app_server,
crate::tui_internal::resume_picker::SessionTarget {
path: Some(root_rollout_path),
thread_id: root_thread_id,
},
)
.await?;
assert!(matches!(control, AppRunControl::Continue));
assert_eq!(app.chat_widget.thread_id(), Some(root_thread_id));
assert_eq!(take_backfill_counts(&requests), (1, expected_reads));
assert_eq!(
app.agent_navigation.get(&child_thread_id),
Some(&AgentPickerThreadEntry {
agent_nickname: Some("worker".to_string()),
agent_role: Some("worker".to_string()),
agent_path: Some("/root/worker".to_string()),
is_running: false,
is_closed: false,
})
);
Box::pin(app.open_agent_picker(&mut app_server)).await;
assert_eq!(take_backfill_counts(&requests), (1, expected_reads + 1));
app_server.shutdown().await?;
proxy.await??;
Ok(())
})
})?
.join()
.expect("session lifecycle request test thread")
}