From 6fbb3665a202f9becf1558ad751d327f311caf81 Mon Sep 17 00:00:00 2001 From: Owen Lin Date: Sun, 2 Nov 2025 12:01:19 -0800 Subject: [PATCH] thread/list and thread/resume --- .../app-server/src/codex_message_processor.rs | 209 +++++++++++++++++- .../app-server/tests/common/mcp_process.rs | 20 ++ codex-rs/app-server/tests/suite/v2/mod.rs | 2 + .../app-server/tests/suite/v2/thread_list.rs | 181 +++++++++++++++ .../tests/suite/v2/thread_resume.rs | 88 ++++++++ 5 files changed, 488 insertions(+), 12 deletions(-) create mode 100644 codex-rs/app-server/tests/suite/v2/thread_list.rs create mode 100644 codex-rs/app-server/tests/suite/v2/thread_resume.rs diff --git a/codex-rs/app-server/src/codex_message_processor.rs b/codex-rs/app-server/src/codex_message_processor.rs index 481ea8de48..7691ea6d5d 100644 --- a/codex-rs/app-server/src/codex_message_processor.rs +++ b/codex-rs/app-server/src/codex_message_processor.rs @@ -62,6 +62,10 @@ use codex_app_server_protocol::SortOrder; use codex_app_server_protocol::Thread; use codex_app_server_protocol::ThreadArchiveParams; use codex_app_server_protocol::ThreadArchiveResponse; +use codex_app_server_protocol::ThreadListParams; +use codex_app_server_protocol::ThreadListResponse; +use codex_app_server_protocol::ThreadResumeParams; +use codex_app_server_protocol::ThreadResumeResponse; use codex_app_server_protocol::ThreadStartParams; use codex_app_server_protocol::ThreadStartResponse; use codex_app_server_protocol::ThreadStartedNotification; @@ -108,6 +112,7 @@ use codex_protocol::items::TurnItem; use codex_protocol::models::ResponseItem; use codex_protocol::protocol::RateLimitSnapshot as CoreRateLimitSnapshot; use codex_protocol::protocol::RolloutItem; +use codex_protocol::protocol::SessionMetaLine; use codex_protocol::protocol::USER_MESSAGE_BEGIN; use codex_protocol::user_input::UserInput as CoreInputItem; use codex_utils_json_to_toml::json_to_toml; @@ -188,22 +193,14 @@ impl CodexMessageProcessor { ClientRequest::ThreadStart { request_id, params } => { self.thread_start(request_id, params).await; } - ClientRequest::ThreadResume { - request_id, - params: _, - } => { - self.send_unimplemented_error(request_id, "thread/resume") - .await; + ClientRequest::ThreadResume { request_id, params } => { + self.thread_resume(request_id, params).await; } ClientRequest::ThreadArchive { request_id, params } => { self.thread_archive(request_id, params).await; } - ClientRequest::ThreadList { - request_id, - params: _, - } => { - self.send_unimplemented_error(request_id, "thread/list") - .await; + ClientRequest::ThreadList { request_id, params } => { + self.thread_list(request_id, params).await; } ClientRequest::ThreadCompact { request_id, @@ -1113,6 +1110,194 @@ impl CodexMessageProcessor { }) } + async fn thread_list(&self, request_id: RequestId, params: ThreadListParams) { + let ThreadListParams { + cursor, + limit, + order: _, + } = params; + + let page_size = limit.unwrap_or(25).max(1) as usize; + + // Decode cursor string to Cursor via serde (Cursor implements Deserialize from string) + let cursor_obj: Option = match cursor { + Some(s) => serde_json::from_str::(&format!("\"{s}\"")).ok(), + None => None, + }; + let cursor_ref = cursor_obj.as_ref(); + + // v2 API does not filter by provider unless specified; include all. + let model_provider_slice: Option<&[String]> = None; + let fallback_provider = self.config.model_provider_id.clone(); + + let page = match RolloutRecorder::list_conversations( + &self.config.codex_home, + page_size, + cursor_ref, + INTERACTIVE_SESSION_SOURCES, + model_provider_slice, + fallback_provider.as_str(), + ) + .await + { + Ok(p) => p, + Err(err) => { + let error = JSONRPCErrorError { + code: INTERNAL_ERROR_CODE, + message: format!("failed to list threads: {err}"), + data: None, + }; + self.outgoing.send_error(request_id, error).await; + return; + } + }; + + let mut data = Vec::new(); + for item in page.items.into_iter() { + if let Some(summary) = + extract_conversation_summary(item.path.clone(), &item.head, &fallback_provider) + { + data.push(Thread { + id: summary.conversation_id.to_string(), + turn: Vec::::new(), + }); + continue; + } + if let Some(first) = item.head.first() { + if let Ok(meta) = serde_json::from_value::(first.clone()) { + data.push(Thread { + id: meta.id.to_string(), + turn: Vec::::new(), + }); + continue; + } + if let Ok(line) = serde_json::from_value::(first.clone()) { + data.push(Thread { + id: line.meta.id.to_string(), + turn: Vec::::new(), + }); + continue; + } + } + } + + // Encode next_cursor as a plain string + let next_cursor = match page.next_cursor { + Some(c) => match serde_json::to_value(&c) { + Ok(serde_json::Value::String(s)) => Some(s), + _ => None, + }, + None => None, + }; + + let response = ThreadListResponse { data, next_cursor }; + self.outgoing.send_response(request_id, response).await; + } + + async fn thread_resume(&self, request_id: RequestId, params: ThreadResumeParams) { + // Convert thread id string into ConversationId. + let conversation_id = match ConversationId::from_string(¶ms.thread.id) { + Ok(id) => id, + Err(err) => { + let error = JSONRPCErrorError { + code: INVALID_REQUEST_ERROR_CODE, + message: format!("invalid thread id: {err}"), + data: None, + }; + self.outgoing.send_error(request_id, error).await; + return; + } + }; + + // Locate the rollout path for the conversation id. + let path = match find_conversation_path_by_id_str( + &self.config.codex_home, + &conversation_id.to_string(), + ) + .await + { + Ok(Some(p)) => p, + Ok(None) => { + let error = JSONRPCErrorError { + code: INVALID_REQUEST_ERROR_CODE, + message: format!("no rollout found for conversation id {conversation_id}"), + data: None, + }; + self.outgoing.send_error(request_id, error).await; + return; + } + Err(err) => { + let error = JSONRPCErrorError { + code: INVALID_REQUEST_ERROR_CODE, + message: format!("failed to locate conversation id {conversation_id}: {err}"), + data: None, + }; + self.outgoing.send_error(request_id, error).await; + return; + } + }; + + // Read initial history from the rollout and resume the conversation. + let fallback_provider = self.config.model_provider_id.as_str(); + let summary = match read_summary_from_rollout(&path, fallback_provider).await { + Ok(s) => s, + Err(err) => { + let error = JSONRPCErrorError { + code: INVALID_REQUEST_ERROR_CODE, + message: format!("failed to load rollout `{}`: {err}", path.display()), + data: None, + }; + self.outgoing.send_error(request_id, error).await; + return; + } + }; + + let initial_history = match RolloutRecorder::get_rollout_history(&summary.path).await { + Ok(initial_history) => initial_history, + Err(err) => { + let error = JSONRPCErrorError { + code: INVALID_REQUEST_ERROR_CODE, + message: format!( + "failed to load rollout `{}` for conversation {conversation_id}: {err}", + summary.path.display() + ), + data: None, + }; + self.outgoing.send_error(request_id, error).await; + return; + } + }; + + match self + .conversation_manager + .resume_conversation_with_history( + self.config.as_ref().clone(), + initial_history, + self.auth_manager.clone(), + ) + .await + { + Ok(_) => { + // Return the same thread id and an empty turn list for now. + let response = ThreadResumeResponse { + thread: Thread { + id: conversation_id.to_string(), + turn: Vec::::new(), + }, + }; + self.outgoing.send_response(request_id, response).await; + } + Err(err) => { + let error = JSONRPCErrorError { + code: INTERNAL_ERROR_CODE, + message: format!("error resuming thread: {err}"), + data: None, + }; + self.outgoing.send_error(request_id, error).await; + } + } + } + async fn get_conversation_summary( &self, request_id: RequestId, diff --git a/codex-rs/app-server/tests/common/mcp_process.rs b/codex-rs/app-server/tests/common/mcp_process.rs index 8806dc1809..7743034a5f 100644 --- a/codex-rs/app-server/tests/common/mcp_process.rs +++ b/codex-rs/app-server/tests/common/mcp_process.rs @@ -32,6 +32,8 @@ use codex_app_server_protocol::SendUserTurnParams; use codex_app_server_protocol::ServerRequest; use codex_app_server_protocol::SetDefaultModelParams; use codex_app_server_protocol::ThreadArchiveParams; +use codex_app_server_protocol::ThreadListParams; +use codex_app_server_protocol::ThreadResumeParams; use codex_app_server_protocol::ThreadStartParams; use codex_app_server_protocol::JSONRPCError; @@ -327,6 +329,24 @@ impl McpProcess { self.send_request("thread/archive", params).await } + /// Send a `thread/list` JSON-RPC request (v2). + pub async fn send_thread_list_request( + &mut self, + params: ThreadListParams, + ) -> anyhow::Result { + let params = Some(serde_json::to_value(params)?); + self.send_request("thread/list", params).await + } + + /// Send a `thread/resume` JSON-RPC request (v2). + pub async fn send_thread_resume_request( + &mut self, + params: ThreadResumeParams, + ) -> anyhow::Result { + let params = Some(serde_json::to_value(params)?); + self.send_request("thread/resume", params).await + } + /// Send a `cancelLoginChatGpt` JSON-RPC request. pub async fn send_cancel_login_chat_gpt_request( &mut self, diff --git a/codex-rs/app-server/tests/suite/v2/mod.rs b/codex-rs/app-server/tests/suite/v2/mod.rs index 7f79cf9718..18b54c743d 100644 --- a/codex-rs/app-server/tests/suite/v2/mod.rs +++ b/codex-rs/app-server/tests/suite/v2/mod.rs @@ -1 +1,3 @@ +mod thread_list; +mod thread_resume; mod thread_start; diff --git a/codex-rs/app-server/tests/suite/v2/thread_list.rs b/codex-rs/app-server/tests/suite/v2/thread_list.rs new file mode 100644 index 0000000000..fd7858b2e7 --- /dev/null +++ b/codex-rs/app-server/tests/suite/v2/thread_list.rs @@ -0,0 +1,181 @@ +use anyhow::Result; +use app_test_support::McpProcess; +use app_test_support::to_response; +use codex_app_server_protocol::JSONRPCResponse; +use codex_app_server_protocol::RequestId; +use codex_app_server_protocol::ThreadListParams; +use codex_app_server_protocol::ThreadListResponse; +use serde_json::json; +use tempfile::TempDir; +use tokio::time::timeout; +use uuid::Uuid; + +const DEFAULT_READ_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(10); + +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn thread_list_basic_empty() -> Result<()> { + let codex_home = TempDir::new()?; + create_minimal_config(codex_home.path())?; + + let mut mcp = McpProcess::new(codex_home.path()).await?; + timeout(DEFAULT_READ_TIMEOUT, mcp.initialize()).await??; + + // List threads in an empty CODEX_HOME; should return an empty page with nextCursor: null. + let list_id = mcp + .send_thread_list_request(ThreadListParams { + cursor: None, + limit: Some(10), + order: None, + }) + .await?; + let list_resp: JSONRPCResponse = timeout( + DEFAULT_READ_TIMEOUT, + mcp.read_stream_until_response_message(RequestId::Integer(list_id)), + ) + .await??; + let ThreadListResponse { data, next_cursor } = to_response::(list_resp)?; + assert!(data.is_empty()); + assert!(next_cursor.is_none()); + + Ok(()) +} + +// Minimal config.toml for listing. +fn create_minimal_config(codex_home: &std::path::Path) -> std::io::Result<()> { + let config_toml = codex_home.join("config.toml"); + std::fs::write( + config_toml, + "model = \"mock-model\"\napproval_policy = \"never\"\n", + ) +} + +fn create_fake_rollout( + codex_home: &std::path::Path, + filename_ts: &str, + meta_rfc3339: &str, +) -> Result { + let uuid = Uuid::new_v4(); + let year = &filename_ts[0..4]; + let month = &filename_ts[5..7]; + let day = &filename_ts[8..10]; + let dir = codex_home.join("sessions").join(year).join(month).join(day); + std::fs::create_dir_all(&dir)?; + + let file_path = dir.join(format!("rollout-{filename_ts}-{uuid}.jsonl")); + let mut lines = Vec::new(); + lines.push( + json!({ + "timestamp": meta_rfc3339, + "type": "session_meta", + "payload": { + "id": uuid, + "timestamp": meta_rfc3339, + "cwd": "/", + "originator": "codex", + "cli_version": "0.0.0", + "instructions": null, + "source": "vscode", + "model_provider": "mock_provider" + } + }) + .to_string(), + ); + lines.push( + json!({ + "timestamp": meta_rfc3339, + "type":"response_item", + "payload": { + "type":"message", + "role":"user", + "content":[{"type":"input_text","text": "Hello"}] + } + }) + .to_string(), + ); + // Add a matching user_message event so the scanner includes this rollout. + lines.push( + json!({ + "timestamp": meta_rfc3339, + "type":"event_msg", + "payload": { + "type":"user_message", + "message":"Hello", + "kind":"plain" + } + }) + .to_string(), + ); + std::fs::write(file_path, lines.join("\n") + "\n")?; + Ok(uuid.to_string()) +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn thread_list_pagination_next_cursor_none_on_last_page() -> Result<()> { + let codex_home = TempDir::new()?; + create_minimal_config(codex_home.path())?; + + // Create three rollouts so we can paginate with limit=2. + let _a = create_fake_rollout( + codex_home.path(), + "2025-01-02T12-00-00", + "2025-01-02T12:00:00Z", + )?; + let _b = create_fake_rollout( + codex_home.path(), + "2025-01-01T13-00-00", + "2025-01-01T13:00:00Z", + )?; + let _c = create_fake_rollout( + codex_home.path(), + "2025-01-01T12-00-00", + "2025-01-01T12:00:00Z", + )?; + + let mut mcp = McpProcess::new(codex_home.path()).await?; + timeout(DEFAULT_READ_TIMEOUT, mcp.initialize()).await??; + + // Page 1: limit 2 → expect next_cursor Some. + let page1_id = mcp + .send_thread_list_request(ThreadListParams { + cursor: None, + limit: Some(2), + order: None, + }) + .await?; + let page1_resp: JSONRPCResponse = timeout( + DEFAULT_READ_TIMEOUT, + mcp.read_stream_until_response_message(RequestId::Integer(page1_id)), + ) + .await??; + let ThreadListResponse { + data: data1, + next_cursor: cursor1, + } = to_response::(page1_resp)?; + assert_eq!(data1.len(), 2); + let cursor1 = cursor1.expect("expected nextCursor on first page"); + + // Page 2: with cursor → expect next_cursor None when no more results. + let page2_id = mcp + .send_thread_list_request(ThreadListParams { + cursor: Some(cursor1), + limit: Some(2), + order: None, + }) + .await?; + let page2_resp: JSONRPCResponse = timeout( + DEFAULT_READ_TIMEOUT, + mcp.read_stream_until_response_message(RequestId::Integer(page2_id)), + ) + .await??; + let ThreadListResponse { + data: data2, + next_cursor: cursor2, + } = to_response::(page2_resp)?; + assert!(data2.len() <= 2); + assert!( + cursor2.is_none(), + "expected nextCursor to be null on last page" + ); + + Ok(()) +} diff --git a/codex-rs/app-server/tests/suite/v2/thread_resume.rs b/codex-rs/app-server/tests/suite/v2/thread_resume.rs new file mode 100644 index 0000000000..bb91ec74bf --- /dev/null +++ b/codex-rs/app-server/tests/suite/v2/thread_resume.rs @@ -0,0 +1,88 @@ +use anyhow::Result; +use app_test_support::McpProcess; +use app_test_support::create_mock_chat_completions_server; +use app_test_support::to_response; +use codex_app_server_protocol::JSONRPCResponse; +use codex_app_server_protocol::RequestId; +use codex_app_server_protocol::ThreadResumeParams; +use codex_app_server_protocol::ThreadResumeResponse; +use codex_app_server_protocol::ThreadStartParams; +use codex_app_server_protocol::ThreadStartResponse; +use tempfile::TempDir; +use tokio::time::timeout; + +const DEFAULT_READ_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(10); + +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn thread_resume_returns_existing_thread() -> Result<()> { + let server = create_mock_chat_completions_server(vec![]).await; + let codex_home = TempDir::new()?; + create_config_toml(codex_home.path(), &server.uri())?; + + let mut mcp = McpProcess::new(codex_home.path()).await?; + timeout(DEFAULT_READ_TIMEOUT, mcp.initialize()).await??; + + // Start a thread. + let start_id = mcp + .send_thread_start_request(ThreadStartParams { + model: Some("o3".to_string()), + model_provider: None, + profile: None, + cwd: None, + approval_policy: None, + sandbox: None, + config: None, + base_instructions: None, + developer_instructions: None, + compact_prompt: None, + include_apply_patch_tool: None, + }) + .await?; + let start_resp: JSONRPCResponse = timeout( + DEFAULT_READ_TIMEOUT, + mcp.read_stream_until_response_message(RequestId::Integer(start_id)), + ) + .await??; + let ThreadStartResponse { thread } = to_response::(start_resp)?; + + // Resume it via v2 API. + let resume_id = mcp + .send_thread_resume_request(ThreadResumeParams { + thread: thread.clone(), + }) + .await?; + let resume_resp: JSONRPCResponse = timeout( + DEFAULT_READ_TIMEOUT, + mcp.read_stream_until_response_message(RequestId::Integer(resume_id)), + ) + .await??; + let ThreadResumeResponse { thread: resumed } = + to_response::(resume_resp)?; + assert_eq!(resumed.id, thread.id); + + Ok(()) +} + +// Helper to create a config.toml pointing at the mock model server. +fn create_config_toml(codex_home: &std::path::Path, server_uri: &str) -> std::io::Result<()> { + let config_toml = codex_home.join("config.toml"); + std::fs::write( + config_toml, + format!( + r#" +model = "mock-model" +approval_policy = "never" +sandbox_mode = "danger-full-access" + +model_provider = "mock_provider" + +[model_providers.mock_provider] +name = "Mock provider for test" +base_url = "{server_uri}/v1" +wire_api = "chat" +request_max_retries = 0 +stream_max_retries = 0 +"# + ), + ) +}