From f968a1327ad39a7786759ea8f1d1c088fe41e91b Mon Sep 17 00:00:00 2001 From: Michael Bolin Date: Wed, 13 Aug 2025 23:12:03 -0700 Subject: [PATCH 1/2] feat: add support for an InterruptConversation request (#2287) This adds `ClientRequest::InterruptConversation`, which effectively maps directly to `Op::Interrupt`. --- * __->__ #2287 * #2286 * #2285 --- .../mcp-server/src/codex_message_processor.rs | 34 +++++++++++++++++++ codex-rs/mcp-server/src/wire_format.rs | 16 ++++++++- 2 files changed, 49 insertions(+), 1 deletion(-) diff --git a/codex-rs/mcp-server/src/codex_message_processor.rs b/codex-rs/mcp-server/src/codex_message_processor.rs index 74b69fbcbf..a13f0f2677 100644 --- a/codex-rs/mcp-server/src/codex_message_processor.rs +++ b/codex-rs/mcp-server/src/codex_message_processor.rs @@ -34,6 +34,8 @@ use crate::wire_format::EXEC_COMMAND_APPROVAL_METHOD; use crate::wire_format::ExecCommandApprovalParams; use crate::wire_format::ExecCommandApprovalResponse; use crate::wire_format::InputItem as WireInputItem; +use crate::wire_format::InterruptConversationParams; +use crate::wire_format::InterruptConversationResponse; use crate::wire_format::NewConversationParams; use crate::wire_format::NewConversationResponse; use crate::wire_format::RemoveConversationListenerParams; @@ -76,6 +78,9 @@ impl CodexMessageProcessor { ClientRequest::SendUserMessage { request_id, params } => { self.send_user_message(request_id, params).await; } + ClientRequest::InterruptConversation { request_id, params } => { + self.interrupt_conversation(request_id, params).await; + } ClientRequest::AddConversationListener { request_id, params } => { self.add_conversation_listener(request_id, params).await; } @@ -164,6 +169,35 @@ impl CodexMessageProcessor { .await; } + async fn interrupt_conversation( + &mut self, + request_id: RequestId, + params: InterruptConversationParams, + ) { + let InterruptConversationParams { conversation_id } = params; + let Ok(conversation) = self + .conversation_manager + .get_conversation(conversation_id.0) + .await + else { + let error = JSONRPCErrorError { + code: INVALID_REQUEST_ERROR_CODE, + message: format!("conversation not found: {conversation_id}"), + data: None, + }; + self.outgoing.send_error(request_id, error).await; + return; + }; + + let _ = conversation.submit(Op::Interrupt).await; + + // Apparently CodexConversation does not send an ack for Op::Interrupt, + // so we can reply to the request right away. + self.outgoing + .send_response(request_id, InterruptConversationResponse {}) + .await; + } + async fn add_conversation_listener( &mut self, request_id: RequestId, diff --git a/codex-rs/mcp-server/src/wire_format.rs b/codex-rs/mcp-server/src/wire_format.rs index 5b223e604d..f2702a1d99 100644 --- a/codex-rs/mcp-server/src/wire_format.rs +++ b/codex-rs/mcp-server/src/wire_format.rs @@ -36,7 +36,11 @@ pub enum ClientRequest { request_id: RequestId, params: SendUserMessageParams, }, - + InterruptConversation { + #[serde(rename = "id")] + request_id: RequestId, + params: InterruptConversationParams, + }, AddConversationListener { #[serde(rename = "id")] request_id: RequestId, @@ -112,6 +116,16 @@ pub struct SendUserMessageParams { pub items: Vec, } +#[derive(Serialize, Deserialize, Debug, Clone, PartialEq)] +#[serde(rename_all = "camelCase")] +pub struct InterruptConversationParams { + pub conversation_id: ConversationId, +} + +#[derive(Serialize, Deserialize, Debug, Clone, PartialEq)] +#[serde(rename_all = "camelCase")] +pub struct InterruptConversationResponse {} + #[derive(Serialize, Deserialize, Debug, Clone, PartialEq)] #[serde(rename_all = "camelCase")] pub struct SendUserMessageResponse {} From 070970af3020492dbb23722893c3e69ae82a6f5e Mon Sep 17 00:00:00 2001 From: Michael Bolin Date: Wed, 13 Aug 2025 23:42:50 -0700 Subject: [PATCH 2/2] exploration: rollback #1602 --- codex-rs/common/src/config_override.rs | 6 +- codex-rs/core/src/codex.rs | 84 ++++--------- codex-rs/core/src/config.rs | 15 --- codex-rs/core/src/protocol.rs | 4 - codex-rs/core/src/rollout.rs | 139 ++-------------------- codex-rs/core/tests/cli_stream.rs | 156 ++++++------------------- 6 files changed, 69 insertions(+), 335 deletions(-) diff --git a/codex-rs/common/src/config_override.rs b/codex-rs/common/src/config_override.rs index c9b18edc7c..610195d6d1 100644 --- a/codex-rs/common/src/config_override.rs +++ b/codex-rs/common/src/config_override.rs @@ -64,11 +64,7 @@ impl CliConfigOverrides { // `-c model=o3` without the quotes. let value: Value = match parse_toml_value(value_str) { Ok(v) => v, - Err(_) => { - // Strip leading/trailing quotes if present - let trimmed = value_str.trim().trim_matches(|c| c == '"' || c == '\''); - Value::String(trimmed.to_string()) - } + Err(_) => Value::String(value_str.to_string()), }; Ok((key.to_string(), value)) diff --git a/codex-rs/core/src/codex.rs b/codex-rs/core/src/codex.rs index ff726a2426..f6a9bd968e 100644 --- a/codex-rs/core/src/codex.rs +++ b/codex-rs/core/src/codex.rs @@ -125,9 +125,6 @@ pub struct CodexSpawnOk { impl Codex { /// Spawn a new [`Codex`] and initialize the session. pub async fn spawn(config: Config, auth: Option) -> CodexResult { - // experimental resume path (undocumented) - let resume_path = config.experimental_resume.clone(); - info!("resume_path: {resume_path:?}"); let (tx_sub, rx_sub) = async_channel::bounded(64); let (tx_event, rx_event) = async_channel::unbounded(); @@ -145,7 +142,6 @@ impl Codex { disable_response_storage: config.disable_response_storage, notify: config.notify.clone(), cwd: config.cwd.clone(), - resume_path: resume_path.clone(), }; let config = Arc::new(config); @@ -354,27 +350,25 @@ impl Session { state.approved_commands.insert(cmd); } - /// Records items to both the rollout and the chat completions/ZDR - /// transcript, if enabled. + /// Records items to both the rollout and the in-memory conversation + /// history, if enabled. async fn record_conversation_items(&self, items: &[ResponseItem]) { debug!("Recording items for conversation: {items:?}"); - self.record_state_snapshot(items).await; - + self.record_rollout_items(items).await; self.state.lock().unwrap().history.record_items(items); } - async fn record_state_snapshot(&self, items: &[ResponseItem]) { - let snapshot = { crate::rollout::SessionStateSnapshot {} }; - + /// Append the given items to the session's rollout transcript (if enabled) + /// and persist them to disk. + async fn record_rollout_items(&self, items: &[ResponseItem]) { + // Clone the recorder outside of the mutex so we don't hold the lock + // across an await point (MutexGuard is not Send). let recorder = { let guard = self.rollout.lock().unwrap(); guard.as_ref().cloned() }; if let Some(rec) = recorder { - if let Err(e) = rec.record_state(snapshot).await { - error!("failed to record rollout state: {e:#}"); - } if let Err(e) = rec.record_items(items).await { error!("failed to record rollout items: {e:#}"); } @@ -709,12 +703,13 @@ impl AgentTask { } async fn submission_loop( - mut session_id: Uuid, + session_id: Uuid, config: Arc, auth: Option, rx_sub: Receiver, tx_event: Sender, ) { + // session_id is provided by the caller (spawn) let mut sess: Option> = None; // shorthand - send an event when there is no active session let send_no_session_event = |sub_id: String| async { @@ -754,11 +749,8 @@ async fn submission_loop( disable_response_storage, notify, cwd, - resume_path, } => { - debug!( - "Configuring session: model={model}; provider={provider:?}; resume={resume_path:?}" - ); + debug!("Configuring session: model={model}; provider={provider:?}"); if !cwd.is_absolute() { let message = format!("cwd is not absolute: {cwd:?}"); error!(message); @@ -771,39 +763,19 @@ async fn submission_loop( } return; } - // Optionally resume an existing rollout. - let mut restored_items: Option> = None; - let rollout_recorder: Option = - if let Some(path) = resume_path.as_ref() { - match RolloutRecorder::resume(path, cwd.clone()).await { - Ok((rec, saved)) => { - session_id = saved.session_id; - if !saved.items.is_empty() { - restored_items = Some(saved.items); - } - Some(rec) - } - Err(e) => { - warn!("failed to resume rollout from {path:?}: {e}"); - None - } - } - } else { + // Attempt to create a RolloutRecorder before moving the + // `user_instructions` value into the Session struct. + let rollout_recorder = match RolloutRecorder::new( + &config, + session_id, + user_instructions.clone(), + ) + .await + { + Ok(r) => Some(r), + Err(e) => { + warn!("failed to initialise rollout recorder: {e}"); None - }; - - let rollout_recorder = match rollout_recorder { - Some(rec) => Some(rec), - None => { - match RolloutRecorder::new(&config, session_id, user_instructions.clone()) - .await - { - Ok(r) => Some(r), - Err(e) => { - warn!("failed to initialise rollout recorder: {e}"); - None - } - } } }; @@ -885,14 +857,6 @@ async fn submission_loop( show_raw_agent_reasoning: config.show_raw_agent_reasoning, })); - // Patch restored state into the newly created session. - if let Some(sess_arc) = &sess { - if restored_items.is_some() { - let mut st = sess_arc.state.lock().unwrap(); - st.history.record_items(restored_items.unwrap().iter()); - } - } - // Gather history metadata for SessionConfiguredEvent. let (history_log_id, history_entry_count) = crate::message_history::history_metadata(&config).await; @@ -961,8 +925,6 @@ async fn submission_loop( } } Op::AddToHistory { text } => { - // TODO: What should we do if we got AddToHistory before ConfigureSession? - // currently, if ConfigureSession has resume path, this history will be ignored let id = session_id; let config = config.clone(); tokio::spawn(async move { diff --git a/codex-rs/core/src/config.rs b/codex-rs/core/src/config.rs index f9c15b9eed..337a16cdd7 100644 --- a/codex-rs/core/src/config.rs +++ b/codex-rs/core/src/config.rs @@ -150,9 +150,6 @@ pub struct Config { /// Base URL for requests to ChatGPT (as opposed to the OpenAI API). pub chatgpt_base_url: String, - /// Experimental rollout resume path (absolute path to .jsonl; undocumented). - pub experimental_resume: Option, - /// Include an experimental plan tool that the model can use to update its current plan and status of each step. pub include_plan_tool: bool, @@ -394,9 +391,6 @@ pub struct ConfigToml { /// Base URL for requests to ChatGPT (as opposed to the OpenAI API). pub chatgpt_base_url: Option, - /// Experimental rollout resume path (absolute path to .jsonl; undocumented). - pub experimental_resume: Option, - /// Experimental path to a file whose contents replace the built-in BASE_INSTRUCTIONS. pub experimental_instructions_file: Option, @@ -593,9 +587,6 @@ impl Config { .as_ref() .map(|info| info.max_output_tokens) }); - - let experimental_resume = cfg.experimental_resume; - // Load base instructions override from a file if specified. If the // path is relative, resolve it against the effective cwd so the // behaviour matches other path-like config values. @@ -606,7 +597,6 @@ impl Config { let file_base_instructions = Self::get_base_instructions(experimental_instructions_path, &resolved_cwd)?; let base_instructions = base_instructions.or(file_base_instructions); - let config = Self { model, model_family, @@ -656,8 +646,6 @@ impl Config { .chatgpt_base_url .or(cfg.chatgpt_base_url) .unwrap_or("https://chatgpt.com/backend-api/".to_string()), - - experimental_resume, include_plan_tool: include_plan_tool.unwrap_or(false), internal_originator: cfg.internal_originator, }; @@ -1020,7 +1008,6 @@ disable_response_storage = true model_reasoning_effort: ReasoningEffort::High, model_reasoning_summary: ReasoningSummary::Detailed, chatgpt_base_url: "https://chatgpt.com/backend-api/".to_string(), - experimental_resume: None, base_instructions: None, include_plan_tool: false, internal_originator: None, @@ -1071,7 +1058,6 @@ disable_response_storage = true model_reasoning_effort: ReasoningEffort::default(), model_reasoning_summary: ReasoningSummary::default(), chatgpt_base_url: "https://chatgpt.com/backend-api/".to_string(), - experimental_resume: None, base_instructions: None, include_plan_tool: false, internal_originator: None, @@ -1137,7 +1123,6 @@ disable_response_storage = true model_reasoning_effort: ReasoningEffort::default(), model_reasoning_summary: ReasoningSummary::default(), chatgpt_base_url: "https://chatgpt.com/backend-api/".to_string(), - experimental_resume: None, base_instructions: None, include_plan_tool: false, internal_originator: None, diff --git a/codex-rs/core/src/protocol.rs b/codex-rs/core/src/protocol.rs index a7e6b1edda..5ab002b7fd 100644 --- a/codex-rs/core/src/protocol.rs +++ b/codex-rs/core/src/protocol.rs @@ -79,10 +79,6 @@ pub enum Op { /// `ConfigureSession` operation so that the business-logic layer can /// operate deterministically. cwd: std::path::PathBuf, - - /// Path to a rollout file to resume from. - #[serde(skip_serializing_if = "Option::is_none")] - resume_path: Option, }, /// Abort current task. diff --git a/codex-rs/core/src/rollout.rs b/codex-rs/core/src/rollout.rs index 0ccd8e891b..6cd6263531 100644 --- a/codex-rs/core/src/rollout.rs +++ b/codex-rs/core/src/rollout.rs @@ -1,13 +1,11 @@ -//! Persist Codex session rollouts (.jsonl) so sessions can be replayed or inspected later. +//! Functionality to persist a Codex conversation rollout to disk. use std::fs::File; use std::fs::{self}; use std::io::Error as IoError; -use std::path::Path; use serde::Deserialize; use serde::Serialize; -use serde_json::Value; use time::OffsetDateTime; use time::format_description::FormatItem; use time::macros::format_description; @@ -15,7 +13,6 @@ use tokio::io::AsyncWriteExt; use tokio::sync::mpsc::Sender; use tokio::sync::mpsc::{self}; use tokio::sync::oneshot; -use tracing::info; use tracing::warn; use uuid::Uuid; @@ -24,6 +21,7 @@ use crate::git_info::GitInfo; use crate::git_info::collect_git_info; use crate::models::ResponseItem; +/// Folder inside `~/.codex` that holds saved rollouts. const SESSIONS_SUBDIR: &str = "sessions"; #[derive(Serialize, Deserialize, Clone, Default)] @@ -41,28 +39,6 @@ struct SessionMetaWithGit { git: Option, } -#[derive(Serialize, Deserialize, Default, Clone)] -pub struct SessionStateSnapshot {} - -#[derive(Serialize, Deserialize, Default, Clone)] -pub struct SavedSession { - pub session: SessionMeta, - #[serde(default)] - pub items: Vec, - #[serde(default)] - pub state: SessionStateSnapshot, - pub session_id: Uuid, -} - -/// Records all [`ResponseItem`]s for a session and flushes them to disk after -/// every update. -/// -/// Rollouts are recorded as JSONL and can be inspected with tools such as: -/// -/// ```ignore -/// $ jq -C . ~/.codex/sessions/rollout-2025-05-07T17-24-21-5973b6c0-94b8-487b-a530-2aeb6098ae0e.jsonl -/// $ fx ~/.codex/sessions/rollout-2025-05-07T17-24-21-5973b6c0-94b8-487b-a530-2aeb6098ae0e.jsonl -/// ``` #[derive(Clone)] pub(crate) struct RolloutRecorder { tx: Sender, @@ -70,7 +46,6 @@ pub(crate) struct RolloutRecorder { enum RolloutCmd { AddItems(Vec), - UpdateState(SessionStateSnapshot), Shutdown { ack: oneshot::Sender<()> }, } @@ -89,6 +64,7 @@ impl RolloutRecorder { timestamp, } = create_log_file(config, uuid)?; + // Build the static session metadata JSON first. let timestamp_format: &[FormatItem] = format_description!( "[year]-[month]-[day]T[hour]:[minute]:[second].[subsecond digits:3]Z" ); @@ -121,8 +97,10 @@ impl RolloutRecorder { Ok(Self { tx }) } + /// Append `items` to the rollout file. pub(crate) async fn record_items(&self, items: &[ResponseItem]) -> std::io::Result<()> { - let mut filtered = Vec::new(); + // Filter out items we should not serialize. + let mut filtered: Vec = Vec::new(); for item in items { match item { // Note that function calls may look a bit strange if they are @@ -131,12 +109,8 @@ impl RolloutRecorder { ResponseItem::Message { .. } | ResponseItem::LocalShellCall { .. } | ResponseItem::FunctionCall { .. } - | ResponseItem::FunctionCallOutput { .. } - | ResponseItem::Reasoning { .. } => filtered.push(item.clone()), - ResponseItem::Other => { - // These should never be serialized. - continue; - } + | ResponseItem::FunctionCallOutput { .. } => filtered.push(item.clone()), + ResponseItem::Reasoning { .. } | ResponseItem::Other => {} } } if filtered.is_empty() { @@ -148,84 +122,6 @@ impl RolloutRecorder { .map_err(|e| IoError::other(format!("failed to queue rollout items: {e}"))) } - pub(crate) async fn record_state(&self, state: SessionStateSnapshot) -> std::io::Result<()> { - self.tx - .send(RolloutCmd::UpdateState(state)) - .await - .map_err(|e| IoError::other(format!("failed to queue rollout state: {e}"))) - } - - pub async fn resume( - path: &Path, - cwd: std::path::PathBuf, - ) -> std::io::Result<(Self, SavedSession)> { - info!("Resuming rollout from {path:?}"); - let text = tokio::fs::read_to_string(path).await?; - let mut lines = text.lines(); - let meta_line = lines - .next() - .ok_or_else(|| IoError::other("empty session file"))?; - let session: SessionMeta = serde_json::from_str(meta_line) - .map_err(|e| IoError::other(format!("failed to parse session meta: {e}")))?; - let mut items = Vec::new(); - let mut state = SessionStateSnapshot::default(); - - for line in lines { - if line.trim().is_empty() { - continue; - } - let v: Value = match serde_json::from_str(line) { - Ok(v) => v, - Err(_) => continue, - }; - if v.get("record_type") - .and_then(|rt| rt.as_str()) - .map(|s| s == "state") - .unwrap_or(false) - { - if let Ok(s) = serde_json::from_value::(v.clone()) { - state = s - } - continue; - } - match serde_json::from_value::(v.clone()) { - Ok(item) => match item { - ResponseItem::Message { .. } - | ResponseItem::LocalShellCall { .. } - | ResponseItem::FunctionCall { .. } - | ResponseItem::FunctionCallOutput { .. } - | ResponseItem::Reasoning { .. } => items.push(item), - ResponseItem::Other => {} - }, - Err(e) => { - warn!("failed to parse item: {v:?}, error: {e}"); - } - } - } - - let saved = SavedSession { - session: session.clone(), - items: items.clone(), - state: state.clone(), - session_id: session.id, - }; - - let file = std::fs::OpenOptions::new() - .append(true) - .read(true) - .open(path)?; - - let (tx, rx) = mpsc::channel::(256); - tokio::task::spawn(rollout_writer( - tokio::fs::File::from_std(file), - rx, - None, - cwd, - )); - info!("Resumed rollout successfully from {path:?}"); - Ok((Self { tx }, saved)) - } - pub async fn shutdown(&self) -> std::io::Result<()> { let (tx_done, rx_done) = oneshot::channel(); match self.tx.send(RolloutCmd::Shutdown { ack: tx_done }).await { @@ -316,28 +212,13 @@ async fn rollout_writer( ResponseItem::Message { .. } | ResponseItem::LocalShellCall { .. } | ResponseItem::FunctionCall { .. } - | ResponseItem::FunctionCallOutput { .. } - | ResponseItem::Reasoning { .. } => { + | ResponseItem::FunctionCallOutput { .. } => { writer.write_line(&item).await?; } - ResponseItem::Other => {} + ResponseItem::Reasoning { .. } | ResponseItem::Other => {} } } } - RolloutCmd::UpdateState(state) => { - #[derive(Serialize)] - struct StateLine<'a> { - record_type: &'static str, - #[serde(flatten)] - state: &'a SessionStateSnapshot, - } - writer - .write_line(&StateLine { - record_type: "state", - state: &state, - }) - .await?; - } RolloutCmd::Shutdown { ack } => { let _ = ack.send(()); } diff --git a/codex-rs/core/tests/cli_stream.rs b/codex-rs/core/tests/cli_stream.rs index c8570577ce..f374f2eb4d 100644 --- a/codex-rs/core/tests/cli_stream.rs +++ b/codex-rs/core/tests/cli_stream.rs @@ -2,6 +2,7 @@ use assert_cmd::Command as AssertCommand; use codex_core::spawn::CODEX_SANDBOX_NETWORK_DISABLED_ENV_VAR; +use serde_json::Value; use std::time::Duration; use std::time::Instant; use tempfile::TempDir; @@ -212,7 +213,6 @@ async fn responses_api_stream_cli() { assert!(stdout.contains("fixture hello")); } -/// End-to-end: create a session (writes rollout), verify the file, then resume and confirm append. #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn integration_creates_and_checks_session_file() { // Honor sandbox network restrictions for CI parity with the other tests. @@ -260,66 +260,45 @@ async fn integration_creates_and_checks_session_file() { String::from_utf8_lossy(&output.stderr) ); - // Wait for sessions dir to appear. + // 5. Sessions are written asynchronously; wait briefly for the directory to appear. let sessions_dir = home.path().join("sessions"); - let dir_deadline = Instant::now() + Duration::from_secs(5); - while !sessions_dir.exists() && Instant::now() < dir_deadline { + let start = Instant::now(); + while !sessions_dir.exists() && start.elapsed() < Duration::from_secs(3) { std::thread::sleep(Duration::from_millis(50)); } - assert!(sessions_dir.exists(), "sessions directory never appeared"); - // Find the session file that contains `marker`. - let deadline = Instant::now() + Duration::from_secs(10); - let mut matching_path: Option = None; - while Instant::now() < deadline && matching_path.is_none() { - for entry in WalkDir::new(&sessions_dir) { - let entry = match entry { - Ok(e) => e, - Err(_) => continue, - }; - if !entry.file_type().is_file() { - continue; - } - if !entry.file_name().to_string_lossy().ends_with(".jsonl") { - continue; - } + // 6. Scan all session files and find the one that contains our marker. + let mut matching_files = vec![]; + for entry in WalkDir::new(&sessions_dir) { + let entry = entry.unwrap(); + if entry.file_type().is_file() && entry.file_name().to_string_lossy().ends_with(".jsonl") { let path = entry.path(); - let Ok(content) = std::fs::read_to_string(path) else { - continue; - }; + let content = std::fs::read_to_string(path).unwrap(); let mut lines = content.lines(); - if lines.next().is_none() { - continue; - } + // Skip SessionMeta (first line) + let _ = lines.next(); for line in lines { - if line.trim().is_empty() { - continue; - } - let item: serde_json::Value = match serde_json::from_str(line) { - Ok(v) => v, - Err(_) => continue, - }; - if item.get("type").and_then(|t| t.as_str()) == Some("message") { - if let Some(c) = item.get("content") { - if c.to_string().contains(&marker) { - matching_path = Some(path.to_path_buf()); + let item: Value = serde_json::from_str(line).unwrap(); + if let Some("message") = item.get("type").and_then(|t| t.as_str()) { + if let Some(content) = item.get("content") { + if content.to_string().contains(&marker) { + matching_files.push(path.to_owned()); break; } } } } } - if matching_path.is_none() { - std::thread::sleep(Duration::from_millis(50)); - } } + assert_eq!( + matching_files.len(), + 1, + "Expected exactly one session file containing the marker, found {}", + matching_files.len() + ); + let path = &matching_files[0]; - let path = match matching_path { - Some(p) => p, - None => panic!("No session file containing the marker was found"), - }; - - // Basic sanity checks on location and metadata. + // 7. Verify directory structure: sessions/YYYY/MM/DD/filename.jsonl let rel = match path.strip_prefix(&sessions_dir) { Ok(r) => r, Err(_) => panic!("session file should live under sessions/"), @@ -348,6 +327,7 @@ async fn integration_creates_and_checks_session_file() { day.len() == 2 && day.chars().all(|c| c.is_ascii_digit()), "Day dir not zero-padded 2-digit numeric: {day}" ); + // Range checks (best-effort; won't fail on leading zeros) if let Ok(m) = month.parse::() { assert!((1..=12).contains(&m), "Month out of range: {m}"); } @@ -355,32 +335,23 @@ async fn integration_creates_and_checks_session_file() { assert!((1..=31).contains(&d), "Day out of range: {d}"); } - let content = - std::fs::read_to_string(&path).unwrap_or_else(|_| panic!("Failed to read session file")); + // 8. Parse SessionMeta line and basic sanity checks. + let content = std::fs::read_to_string(path).unwrap(); let mut lines = content.lines(); - let meta_line = lines - .next() - .ok_or("missing session meta line") - .unwrap_or_else(|_| panic!("missing session meta line")); - let meta: serde_json::Value = serde_json::from_str(meta_line) - .unwrap_or_else(|_| panic!("Failed to parse session meta line as JSON")); + let meta: Value = serde_json::from_str(lines.next().unwrap()).unwrap(); assert!(meta.get("id").is_some(), "SessionMeta missing id"); assert!( meta.get("timestamp").is_some(), "SessionMeta missing timestamp" ); + // 9. Confirm at least one message contains the marker. let mut found_message = false; for line in lines { - if line.trim().is_empty() { - continue; - } - let Ok(item) = serde_json::from_str::(line) else { - continue; - }; - if item.get("type").and_then(|t| t.as_str()) == Some("message") { - if let Some(c) = item.get("content") { - if c.to_string().contains(&marker) { + let item: Value = serde_json::from_str(line).unwrap(); + if item.get("type").map(|t| t == "message").unwrap_or(false) { + if let Some(content) = item.get("content") { + if content.to_string().contains(&marker) { found_message = true; break; } @@ -391,64 +362,7 @@ async fn integration_creates_and_checks_session_file() { found_message, "No message found in session file containing the marker" ); - - // Second run: resume and append. - let orig_len = content.lines().count(); - let marker2 = format!("integration-resume-{}", Uuid::new_v4()); - let prompt2 = format!("echo {marker2}"); - // Cross‑platform safe resume override. On Windows, backslashes in a TOML string must be escaped - // or the parse will fail and the raw literal (including quotes) may be preserved all the way down - // to Config, which in turn breaks resume because the path is invalid. Normalize to forward slashes - // to sidestep the issue. - let resume_path_str = path.to_string_lossy().replace('\\', "/"); - let resume_override = format!("experimental_resume=\"{resume_path_str}\""); - let mut cmd2 = AssertCommand::new("cargo"); - cmd2.arg("run") - .arg("-p") - .arg("codex-cli") - .arg("--quiet") - .arg("--") - .arg("exec") - .arg("--skip-git-repo-check") - .arg("-c") - .arg(&resume_override) - .arg("-C") - .arg(env!("CARGO_MANIFEST_DIR")) - .arg(&prompt2); - cmd2.env("CODEX_HOME", home.path()) - .env("OPENAI_API_KEY", "dummy") - .env("CODEX_RS_SSE_FIXTURE", &fixture) - .env("OPENAI_BASE_URL", "http://unused.local"); - - let output2 = cmd2.output().unwrap(); - assert!(output2.status.success(), "resume codex-cli run failed"); - - // The rollout writer runs on a background async task; give it a moment to flush. - let mut new_len = orig_len; - let deadline = Instant::now() + Duration::from_secs(5); - let mut content2 = String::new(); - while Instant::now() < deadline { - if let Ok(c) = std::fs::read_to_string(&path) { - let count = c.lines().count(); - if count > orig_len { - content2 = c; - new_len = count; - break; - } - } - std::thread::sleep(Duration::from_millis(50)); - } - if content2.is_empty() { - // last attempt - content2 = std::fs::read_to_string(&path).unwrap(); - new_len = content2.lines().count(); - } - assert!(new_len > orig_len, "rollout file did not grow after resume"); - assert!(content2.contains(&marker), "rollout lost original marker"); - assert!( - content2.contains(&marker2), - "rollout missing resumed marker" - ); + // No resume on second run; resume feature removed. } /// Integration test to verify git info is collected and recorded in session files.