diff --git a/pr-body.md b/pr-body.md new file mode 100644 index 0000000..8c85d25 --- /dev/null +++ b/pr-body.md @@ -0,0 +1,35 @@ +## Problem + +During normal Agent usage, the flow silently stops mid-conversation. The AI says "now I'll read the files" and then never continues. No tool call, no more text, no error. It just freezes. + +## Root Causes + +### 1. Actor not emitting end-stream on channel close (actor.rs) +When the command channel closes or receives Abort, the actor only called handle.cancel() which sets the cancellation token but never emits an end-stream frame. The SSE stream hangs indefinitely. + +### 2. lifecycle encoding failure skips close_output() (lifecycle.rs) +Both fail() and cancel() used the ? operator on encode_error_end_stream(), meaning if serialization fails, close_output() is never called. The stream stays open silently. + +### 3. Context/blob sync timeouts too aggressive (context_sync.rs, blob_sync.rs) +Hardcoded 15-second timeouts kill sessions silently on slow networks or with proxy configurations. + +### 4. Bash tool not recognized (dispatch/mod.rs) +Cursor IDE sends "Bash" tool calls but the server only recognized "Shell", causing "unsupported tool: Bash" errors. + +## Changes + +- actor.rs: Use lifecycle::cancel() instead of handle.cancel() +- lifecycle.rs: Use match instead of ? with fallback end-stream frame +- sessions.rs: Document correct Notify pattern +- context_sync.rs: 15s to 60s timeout +- blob_sync.rs: 15s to 60s timeout (SET and GET) +- dispatch/mod.rs: Add Bash tool alias for Shell +- openai_chat.rs: Add detailed debug logging +- router.rs: Add detailed debug logging + +## Testing + +- All existing tests pass +- cargo check compiles cleanly +- Both debug and release builds succeed +- Tested with Cursor BYOK desktop app \ No newline at end of file diff --git a/server/src/cursor/actor.rs b/server/src/cursor/actor.rs index 3df974d..cb03b5c 100644 --- a/server/src/cursor/actor.rs +++ b/server/src/cursor/actor.rs @@ -58,14 +58,14 @@ impl CursorActor { let command = match receiver.recv().await { Some(command) => command, None => { - handle.cancel(); + lifecycle::cancel(&handle).ok(); break; } }; match command { CursorCommand::Abort => { handle.mark_conversation_cancelled(); - handle.cancel(); + lifecycle::cancel(&handle).ok(); } CursorCommand::Finished => { break; diff --git a/server/src/cursor/blob_sync.rs b/server/src/cursor/blob_sync.rs index be4c9dc..ce8acd8 100644 --- a/server/src/cursor/blob_sync.rs +++ b/server/src/cursor/blob_sync.rs @@ -129,7 +129,7 @@ impl BlobSynchronizer { let result = tokio::select! { result = receiver => result.map_err(|_| Error::Protocol("KV SET response channel closed".into()))?, _ = cancellation.cancelled() => Err(Error::Cancelled), - _ = tokio::time::sleep(Duration::from_secs(15)) => Err(Error::Protocol(format!("KV SET timed out: {}", blob_id.to_base64()))), + _ = tokio::time::sleep(Duration::from_secs(60)) => Err(Error::Protocol(format!("KV SET timed out: {}", blob_id.to_base64()))), }; if result.is_err() { self.inner.set_requests.lock().await.remove(&id); @@ -182,7 +182,7 @@ impl BlobSynchronizer { let result = tokio::select! { result = receiver => result.map_err(|_| Error::Protocol("KV GET response channel closed".into()))?, _ = cancellation.cancelled() => Err(Error::Cancelled), - _ = tokio::time::sleep(Duration::from_secs(15)) => Err(Error::Protocol(format!("KV GET timed out: {}", blob_id.to_base64()))), + _ = tokio::time::sleep(Duration::from_secs(60)) => Err(Error::Protocol(format!("KV GET timed out: {}", blob_id.to_base64()))), }; if result.is_err() { self.inner.get_requests.lock().await.remove(&id); diff --git a/server/src/cursor/context_sync.rs b/server/src/cursor/context_sync.rs index 88caa00..b373e9d 100644 --- a/server/src/cursor/context_sync.rs +++ b/server/src/cursor/context_sync.rs @@ -84,7 +84,7 @@ impl RequestContextSynchronizer { let result = tokio::select! { result = receiver => result.map_err(|_| Error::Protocol("request context response channel closed".into()))?, _ = cancellation.cancelled() => Err(Error::Cancelled), - _ = tokio::time::sleep(Duration::from_secs(15)) => Err(Error::Protocol("request context timed out".into())), + _ = tokio::time::sleep(Duration::from_secs(60)) => Err(Error::Protocol("request context timed out".into())), }; if result.is_err() { self.pending.lock().await.take(); diff --git a/server/src/cursor/lifecycle.rs b/server/src/cursor/lifecycle.rs index 6195a41..a9d65b8 100644 --- a/server/src/cursor/lifecycle.rs +++ b/server/src/cursor/lifecycle.rs @@ -32,17 +32,25 @@ pub fn fail(handle: &CursorSessionHandle, error: &Error) -> Result<()> { | Error::Encode(_) | Error::Io(_) => plain_error(ConnectCode::Internal, error), }; - handle.emit_frame(encode_error_end_stream(&stream_error)?); + // Always close the output even if encoding fails, to prevent silent hangs. + match encode_error_end_stream(&stream_error) { + Ok(frame) => handle.emit_frame(frame), + Err(_) => handle.emit_frame(encode_end_stream()), + } handle.close_output(); Ok(()) } pub fn cancel(handle: &CursorSessionHandle) -> Result<()> { - handle.emit_frame(encode_error_end_stream(&ConnectStreamError { + // Always close the output even if encoding fails, to prevent silent hangs. + match encode_error_end_stream(&ConnectStreamError { code: ConnectCode::Canceled, message: "run was cancelled".into(), details: Vec::new(), - })?); + }) { + Ok(frame) => handle.emit_frame(frame), + Err(_) => handle.emit_frame(encode_end_stream()), + } handle.close_output(); Ok(()) } diff --git a/server/src/cursor/sessions.rs b/server/src/cursor/sessions.rs index 5e3b456..6b4b8a8 100644 --- a/server/src/cursor/sessions.rs +++ b/server/src/cursor/sessions.rs @@ -313,6 +313,8 @@ impl CursorSessionRegistry { pub(crate) async fn wait_route(&self, request_id: &str) -> CursorRoute { loop { + // Create the notification future BEFORE checking state to avoid + // a race where a notification fires between state check and await. let changed = self.inner.route_changed.notified(); if self.inner.runs.lock().await.contains_key(request_id) { return CursorRoute::Local; diff --git a/server/src/cursor/tools/dispatch/mod.rs b/server/src/cursor/tools/dispatch/mod.rs index 39957ba..2e6b3ce 100644 --- a/server/src/cursor/tools/dispatch/mod.rs +++ b/server/src/cursor/tools/dispatch/mod.rs @@ -54,7 +54,7 @@ pub(super) async fn start( let call = normalized_call.as_ref().unwrap_or(call); match normalized(&call.name).as_str() { - "shell" | "read" | "delete" | "grep" | "glob" | "readlints" | "task" | "callmcptool" + "shell" | "bash" | "read" | "delete" | "grep" | "glob" | "readlints" | "task" | "callmcptool" | "fetchmcpresource" | "getmcptools" => exec::start(runtime, call, context).await, "write" | "strreplace" | "editnotebook" => edit::start(runtime, call, context).await, "askquestion" | "websearch" | "webfetch" | "switchmode" | "createplan" diff --git a/server/src/provider/openai_chat.rs b/server/src/provider/openai_chat.rs index 75827e7..996fa8e 100644 --- a/server/src/provider/openai_chat.rs +++ b/server/src/provider/openai_chat.rs @@ -15,6 +15,29 @@ use crate::{ Error, Result, }; +macro_rules! dbg_log { + ($loc:expr, $msg:expr, $data:expr) => { + { + let _ = (|| -> std::result::Result<(), Box> { + use std::io::Write; + let payload = serde_json::json!({ + "sessionId": "216d24", + "location": $loc, + "message": $msg, + "data": $data, + "timestamp": std::time::SystemTime::now().duration_since(std::time::UNIX_EPOCH)?.as_millis() as u64, + }); + let mut f = std::fs::OpenOptions::new().create(true).append(true).open( + concat!(env!("CARGO_MANIFEST_DIR"), "/../.cursor/debug-216d24.log") + )?; + writeln!(f, "{}", payload)?; + f.flush()?; + Ok(()) + })(); + } + }; +} + use super::{ apply_openai_prompt_cache_key, merge_extra_params, recorder::recorded_headers, CallRecorder, FinishReason, ModelEvent, Provider, ProviderStream, @@ -61,6 +84,13 @@ impl Provider for OpenAiChatProvider { let recorder = self.recorder.clone(); Box::pin(try_stream! { let ModelInvocation { call_id, request, .. } = invocation; + dbg_log!("openai_chat.rs:stream", "Provider stream started", serde_json::json!({ + "model": request.model.model_id, + "url": config.request_url, + "call_id": call_id.clone(), + "history_len": request.history.len(), + "tools_count": request.prompt.tools.len() + })); let messages = openai_chat_messages(&request.prompt.instructions, &request.history)?; let mut body = json!({ "model": request.model.model_id, @@ -85,7 +115,22 @@ impl Provider for OpenAiChatProvider { _ = cancellation.cancelled() => return, response = request => response, }; - let response = response?; + let response = match response { + Ok(r) => { + dbg_log!("openai_chat.rs:stream", "HTTP response received", serde_json::json!({ + "status": r.status().as_u16(), + "url": config.request_url + })); + r + } + Err(e) => { + dbg_log!("openai_chat.rs:stream", "HTTP request FAILED", serde_json::json!({ + "error": e.to_string(), + "url": config.request_url + })); + Err(Error::from(e))? + } + }; if let Some(recorder) = &recorder { recorder.response_headers(response.status().as_u16()).await?; } @@ -118,15 +163,35 @@ impl Provider for OpenAiChatProvider { let mut final_usage = None; let mut finish = None; let mut saw_done_marker = false; + let mut loop_iteration: u64 = 0; loop { + loop_iteration += 1; let event = tokio::select! { _ = cancellation.cancelled() => { + dbg_log!("openai_chat.rs:stream", "Stream cancelled by token", serde_json::json!({ + "iteration": loop_iteration, + "saw_done_marker": saw_done_marker, + "tool_count": tools.len() + })); return; } event = source.next() => event, }; - let Some(event) = event else { break }; - let event = event.map_err(|error| Error::Provider(format!("OpenAI Chat SSE: {error}")))?; + let Some(event) = event else { + dbg_log!("openai_chat.rs:stream", "SSE stream ended (None)", serde_json::json!({ + "iteration": loop_iteration, + "saw_done_marker": saw_done_marker + })); + break; + }; + let event = event.map_err(|error| { + let err_msg = error.to_string(); + dbg_log!("openai_chat.rs:stream", "SSE event error", serde_json::json!({ + "iteration": loop_iteration, + "error": err_msg.clone() + })); + Error::Provider(format!("OpenAI Chat SSE: {err_msg}")) + })?; if event.data == "[DONE]" { saw_done_marker = true; break; } let value: Value = serde_json::from_str(&event.data)?; if let Some(usage) = value.get("usage").filter(|value| !value.is_null()) { diff --git a/server/src/provider/router.rs b/server/src/provider/router.rs index eeec35f..5a29b0d 100644 --- a/server/src/provider/router.rs +++ b/server/src/provider/router.rs @@ -11,6 +11,29 @@ use crate::{ Error, Result, }; +macro_rules! dbg_log { + ($loc:expr, $msg:expr, $data:expr) => { + { + let _ = (|| -> std::result::Result<(), Box> { + use std::io::Write; + let payload = serde_json::json!({ + "sessionId": "216d24", + "location": $loc, + "message": $msg, + "data": $data, + "timestamp": std::time::SystemTime::now().duration_since(std::time::UNIX_EPOCH)?.as_millis() as u64, + }); + let mut f = std::fs::OpenOptions::new().create(true).append(true).open( + concat!(env!("CARGO_MANIFEST_DIR"), "/../.cursor/debug-216d24.log") + )?; + writeln!(f, "{}", payload)?; + f.flush()?; + Ok(()) + })(); + } + }; +} + use super::{ normalize::NormalizedProvider, AnthropicProvider, CallRecorder, OpenAiChatProvider, OpenAiResponsesProvider, Provider, ProviderStream, @@ -90,23 +113,73 @@ impl Provider for ProviderRouter { let provider = build_observed(&config, recorder.clone(), client)?; let stream_cancellation = cancellation.clone(); let mut stream = provider.stream(invocation, cancellation); + let stream_started = std::time::Instant::now(); + dbg_log!("router.rs:stream", "provider stream created", serde_json::json!({ + "model": selected, + "provider_type": format!("{:?}", provider_type), + "url": config.request_url, + "timeout_ms": config.request_timeout.as_millis() as u64, + })); + let mut last_event_time = std::time::Instant::now(); + let mut event_count: u64 = 0; while let Some(event) = stream.next().await { + let now = std::time::Instant::now(); + let gap_ms = now.duration_since(last_event_time).as_millis() as u64; + let elapsed_ms = now.duration_since(stream_started).as_millis() as u64; + event_count += 1; match event { Ok(event) => { + let event_name = match &event { + super::ModelEvent::Start { .. } => "Start", + super::ModelEvent::TextStart => "TextStart", + super::ModelEvent::TextDelta(d) => { dbg_log!("router.rs:stream", "TextDelta", serde_json::json!({"gap_ms": gap_ms, "elapsed_ms": elapsed_ms, "delta_len": d.len(), "event_count": event_count})); "TextDelta" }, + super::ModelEvent::TextEnd => "TextEnd", + super::ModelEvent::ThinkingStart => "ThinkingStart", + super::ModelEvent::ThinkingDelta(d) => { dbg_log!("router.rs:stream", "ThinkingDelta", serde_json::json!({"gap_ms": gap_ms, "elapsed_ms": elapsed_ms, "delta_len": d.len(), "event_count": event_count})); "ThinkingDelta" }, + super::ModelEvent::ThinkingEnd => "ThinkingEnd", + super::ModelEvent::ToolCallStart { name, .. } => { dbg_log!("router.rs:stream", "ToolCallStart", serde_json::json!({"name": name, "elapsed_ms": elapsed_ms})); "ToolCallStart" }, + super::ModelEvent::ToolCallArgumentsDelta { .. } => "ToolCallArgsDelta", + super::ModelEvent::ToolCallEnd { .. } => "ToolCallEnd", + super::ModelEvent::ProviderReplayState(_) => "ReplayState", + super::ModelEvent::Usage(_) => "Usage", + super::ModelEvent::Done(reason) => { dbg_log!("router.rs:stream", "Done", serde_json::json!({"reason": format!("{:?}", reason), "elapsed_ms": elapsed_ms, "event_count": event_count})); "Done" }, + }; + if gap_ms > 5000 { + dbg_log!("router.rs:stream", "SLOW GAP detected between events", serde_json::json!({ + "gap_ms": gap_ms, + "elapsed_ms": elapsed_ms, + "event": event_name, + "event_count": event_count, + })); + } recorder.event(&event).await?; + last_event_time = now; yield event; } Err(error) => { + dbg_log!("router.rs:stream", "PROVIDER ERROR", serde_json::json!({ + "error": error.to_string(), + "elapsed_ms": elapsed_ms, + "gap_ms": gap_ms, + "event_count": event_count, + })); recorder.failed(&error).await?; Err(error)?; } } } if !recorder.is_finished() { + let elapsed_ms = stream_started.elapsed().as_millis() as u64; if stream_cancellation.is_cancelled() { + dbg_log!("router.rs:stream", "stream ended: cancelled", serde_json::json!({"elapsed_ms": elapsed_ms, "event_count": event_count})); recorder.cancelled().await?; } else { let error = Error::Provider("provider stream ended without Done".into()); + dbg_log!("router.rs:stream", "stream ended: WITHOUT DONE (fatal)", serde_json::json!({ + "elapsed_ms": elapsed_ms, + "event_count": event_count, + "error": "provider stream ended without Done", + })); recorder.failed(&error).await?; Err(error)?; }