diff --git a/src/agent/thread_ops.rs b/src/agent/thread_ops.rs index bc2fc81f..879862c2 100644 --- a/src/agent/thread_ops.rs +++ b/src/agent/thread_ops.rs @@ -1016,8 +1016,19 @@ impl Agent { // Reset thread state to processing { let mut sess = session.lock().await; - if let Some(thread) = sess.threads.get_mut(&thread_id) { - thread.state = ThreadState::Processing; + match sess.threads.get_mut(&thread_id) { + Some(thread) => { + thread.state = ThreadState::Processing; + } + None => { + tracing::error!( + %thread_id, + "Thread disappeared while setting state to Processing during approval" + ); + return Ok(SubmissionResult::error( + "Internal error: thread no longer exists", + )); + } } } @@ -1101,13 +1112,21 @@ impl Agent { // Record sanitized result in thread { let mut sess = session.lock().await; - if let Some(thread) = sess.threads.get_mut(&thread_id) - && let Some(turn) = thread.last_turn_mut() - { - if is_tool_error { - turn.record_tool_error(result_content.clone()); - } else { - turn.record_tool_result(serde_json::json!(result_content)); + match sess.threads.get_mut(&thread_id) { + Some(thread) => { + if let Some(turn) = thread.last_turn_mut() { + if is_tool_error { + turn.record_tool_error(result_content.clone()); + } else { + turn.record_tool_result(serde_json::json!(result_content)); + } + } + } + None => { + tracing::error!( + %thread_id, + "Thread disappeared while recording tool result during approval" + ); } } } @@ -1355,13 +1374,22 @@ impl Agent { // Record sanitized result in thread { let mut sess = session.lock().await; - if let Some(thread) = sess.threads.get_mut(&thread_id) - && let Some(turn) = thread.last_turn_mut() - { - if is_deferred_error { - turn.record_tool_error(deferred_content.clone()); - } else { - turn.record_tool_result(serde_json::json!(deferred_content)); + match sess.threads.get_mut(&thread_id) { + Some(thread) => { + if let Some(turn) = thread.last_turn_mut() { + if is_deferred_error { + turn.record_tool_error(deferred_content.clone()); + } else { + turn.record_tool_result(serde_json::json!(deferred_content)); + } + } + } + None => { + tracing::error!( + %thread_id, + tool_name = %tc.name, + "Thread disappeared while recording deferred tool result during approval" + ); } } } @@ -1414,8 +1442,19 @@ impl Agent { { let mut sess = session.lock().await; - if let Some(thread) = sess.threads.get_mut(&thread_id) { - thread.await_approval(new_pending); + match sess.threads.get_mut(&thread_id) { + Some(thread) => { + thread.await_approval(new_pending); + } + None => { + tracing::error!( + %thread_id, + "Thread disappeared while setting up deferred tool approval" + ); + return Ok(SubmissionResult::error( + "Internal error: thread no longer exists", + )); + } } } @@ -1547,17 +1586,28 @@ impl Agent { ); { let mut sess = session.lock().await; - if let Some(thread) = sess.threads.get_mut(&thread_id) { - thread.clear_pending_approval(); - thread.complete_turn(&rejection); - // User message already persisted at turn start; save rejection response - self.persist_assistant_response( - thread_id, - &message.channel, - &message.user_id, - &rejection, - ) - .await; + match sess.threads.get_mut(&thread_id) { + Some(thread) => { + thread.clear_pending_approval(); + thread.complete_turn(&rejection); + // User message already persisted at turn start; save rejection response + self.persist_assistant_response( + thread_id, + &message.channel, + &message.user_id, + &rejection, + ) + .await; + } + None => { + tracing::error!( + %thread_id, + "Thread disappeared during approval rejection" + ); + return Ok(SubmissionResult::error( + "Internal error: thread no longer exists", + )); + } } } @@ -2099,6 +2149,35 @@ mod tests { } } + #[test] + fn test_approval_on_missing_thread_should_error() { + // Regression for #1487: process_approval() must return an error when + // the thread disappears during approval processing, rather than + // silently succeeding. + use crate::agent::session::{Session, Thread, ThreadState}; + use uuid::Uuid; + + let thread_id = Uuid::new_v4(); + let session_id = Uuid::new_v4(); + let mut session = Session::new("test-user"); + + // Thread doesn't exist in session — simulates disappearance after + // lock re-acquisition. + assert!(!session.threads.contains_key(&thread_id)); + + // The match arms we added should detect this and return an error. + // Verify the lookup returns None — the code path must handle it. + assert!(session.threads.get_mut(&thread_id).is_none()); + + // Also verify a thread that existed and was removed is caught. + let mut thread = Thread::with_id(thread_id, session_id); + thread.start_turn("pending approval"); + thread.state = ThreadState::AwaitingApproval; + session.threads.insert(thread_id, thread); + session.threads.remove(&thread_id); + assert!(session.threads.get_mut(&thread_id).is_none()); + } + #[test] fn test_queue_cap_rejects_at_capacity() { use crate::agent::session::{MAX_PENDING_MESSAGES, Thread, ThreadState};