Compare commits

...
Author SHA1 Message Date
leookun 59379f1f72 fix: repair stream lifecycle reliability 2026-08-28 13:00:39 +08:00
leookun 370f120b8a docs: update README to reflect recent changes and improve clarity 2026-08-28 12:50:01 +08:00
leookun 1ab86c1851 docs: update image references in README for consistency 2026-08-28 12:44:18 +08:00
leokun d015d1c33b docs: sync README from source repo 2026-08-28 12:36:30 +08:00
leokun 1fcb78ad95 Merge pull request #366 from joel01-dev/fix/stream-reliability-v2 2026-08-28 12:14:10 +08:00
joel01-dev b2610c562e fix: restore upstream openai_chat.rs, add only debug logging 2026-08-28 01:12:31 -03:00
joel01-dev 318545cebd fix: restore upstream dispatch/mod.rs, add only Bash alias 2026-08-28 01:07:41 -03:00
joel01-dev 2369d101ec docs: add PR description 2026-08-28 00:48:43 -03:00
joel01-dev 80604bdf97 fix: prevent silent SSE stream hangs, add Bash alias, debug logging
- Actor: use lifecycle::cancel() instead of handle.cancel() on channel
  close and Abort, ensuring end-stream frame is always emitted
- Lifecycle: use match instead of ? on encode_error_end_stream with
  fallback end-stream frame to prevent silent hangs
- Dispatch: add Bash tool alias for Shell
- Context/Blob sync: increase timeouts from 15s to 60s
- Providers: add detailed debug logging to openai_chat and router
2026-08-28 00:47:53 -03:00
12 changed files with 418 additions and 21 deletions
+214
View File
@@ -0,0 +1,214 @@
<div align="center">
# cursor-byok
cursor-byok 是一个运行在本机的 Cursor 模型网关,帮助你在 Cursor 中使用自己配置的模型服务。
[English README](./README.md) · [使用指南](https://docs.leokun.cn) · [下载](https://github.com/leookun/cursor-byok/releases/latest) · [提交问题](https://github.com/leookun/cursor-byok/issues)
[![Release](https://img.shields.io/github/v/release/leookun/cursor-byok?style=flat-square)](https://github.com/leookun/cursor-byok/releases/latest)
[![Downloads](https://img.shields.io/github/downloads/leookun/cursor-byok/total?style=flat-square)](https://github.com/leookun/cursor-byok/releases)
[![License](https://img.shields.io/github/license/leookun/cursor-byok?style=flat-square)](./LICENSE)
[![Platforms](https://img.shields.io/badge/platform-macOS%20%7C%20Windows%20%7C%20Linux-lightgrey?style=flat-square)](https://github.com/leookun/cursor-byok/releases/latest)
</div>
![将 cursor-byok 连接到多种模型 API](./images/en-brand-1.png)
![cursor-byok 控制面板](./images/en-home-1.png)
## 项目简介
cursor-byok 是一个开源的本地模型网关。它在你的设备上运行服务,接收 Cursor 发出的 Agent 请求,将请求转发到你配置的模型服务,并尽可能保留 Cursor Agent 的工具调用、Skills、MCP 和多轮对话能力。
你可以连接兼容 OpenAI 或 Anthropic 协议的服务,自定义服务地址、模型 ID、API Key 和请求参数,也可以使用 Cursor 平台默认选项之外的模型通道。
> [!IMPORTANT]
> cursor-byok 免费且开源,但你连接的模型服务商可能会按用量收费。本项目是独立项目,与 Cursor 或其开发者没有关联,也未获得其认可。
## 主要功能
- **自定义模型通道**:配置自己的 API 地址、凭据和模型 ID。
- **多种 API 协议**:支持 OpenAI Responses API、OpenAI Chat Completions API 和 Anthropic Messages API 兼容服务。
- **模型管理**:添加、复制、编辑、排序模型配置,并批量测试连接。
- **连接性能测试**:查看首字延迟、生成速度、总耗时和原始服务商响应。
- **Agent 工作流**:继续使用工具调用、Skills、MCP 和多轮对话。
- **会话指标**:查看 Token 用量、缓存命中率、对话轮次和估算价值。
- **TAB 补全服务**:在公益服务、官方直连和自定义服务之间选择连接方式。
- **跨平台运行**:支持 macOS、Windows 和 Linux。
## 快速开始
1. 从 [GitHub Releases](https://github.com/leookun/cursor-byok/releases/latest) 下载适合你操作系统的最新版本。
2. 启动 cursor-byok,打开 **Cursor 配置**,按提示初始化本地 CA(证书颁发机构)。
3. 在模型设置中添加模型,填写服务地址、API Key 和模型名称,然后保存并运行 **测试**。
4. 确认测试通过后,保持 cursor-byok 运行。
5. 打开或重启 Cursor,在模型列表中选择已配置的模型,开始使用 Agent。
完整的安装步骤、配置说明和故障处理,请参阅[中文使用指南](https://docs.leokun.cn/zh/docs)。
> [!TIP]
> 首次完成配置后,建议完全退出并重新启动 Cursor,再新开一个对话。使用自定义模型时,请在模型列表中手动选择该模型,不要选择 **Auto**。
## 模型配置
每个模型配置都是独立的上游通道,可以单独设置服务商、协议、凭据和生成参数。
![模型设置页面](./images/en-model-1.png)
### 类型与协议选择
| 模型系列 | 模型类型 | 请求协议 |
| --- | --- | --- |
| Claude 系列 | **Anthropic** | Messages API |
| GPT / OpenAI 系列 | **OpenAI** | **Responses API** |
| 其他模型 | **OpenAI** | **Chat Completions API** |
GPT 系列建议使用 **Responses API**。如果使用 Chat Completions,可能无法保留提示词缓存,导致速度变慢和费用增加。
### 常用字段
- **模型类型**:选择 OpenAI 或 Anthropic,决定上游接口格式。
- **请求协议**:OpenAI 类型需要继续选择 Responses API 或 Chat Completions API。
- **服务器地址**:可以填写服务商基础地址,让应用按协议追加标准端点,也可以填写完整请求 URL 并原样使用。
- **API Key**:填写上游服务要求的访问密钥。密钥保存在本机,用于发送模型请求。
- **模型名称**:填写服务商接口接受的模型标识,也可以使用 **获取模型** 读取服务商返回的模型列表。
- **显示名称**:Cursor 模型列表中显示的名称,不会改变发送给上游的模型标识。
- **备注**:显示在 Cursor 的模型说明中。
还可以根据模型能力设置上下文窗口 Token、最大输出 Token、推理或思考强度、自定义 Headers,以及 OpenAI 或 Anthropic 的额外参数。自定义 Headers 和额外参数必须是 JSON 对象,只应填写服务商明确支持的字段。
保存配置后运行 **测试**,确认地址、协议、API Key、模型标识和流式响应都正常,再在 Cursor 中使用该模型。
## TAB 补全服务
Cursor 的 Tab 补全由独立的 TAB 服务处理,不经过模型通道。你可以在 **系统设置 → TAB 设置** 中选择以下模式:
- **使用公益服务(默认)**:使用项目作者部署的公共服务,无需额外配置。
- **直连**:直接连接当前 Cursor 账号对应的官方 TAB 服务,适合账号拥有官方额度的情况。
- **自定义**:自行部署 [`cursor-tab-server`](https://github.com/leookun/cursor-byok/tree/archive/v0.0.49/cursor-tab-server),然后填写 TAB 服务地址。
修改 TAB 设置后,建议重启 Cursor 并新开一个对话,确保新的连接方式生效。
## 与官方账号并存
新版设计支持 cursor-byok 与 Cursor 官方服务并存:
- 直接在 Cursor 中登录自己的账号。如果之前使用旧版生成的 fake 账户,请先退出该账户,再登录自己的账号。
- 账号拥有官方额度时,官方模型和本地模型可以随时切换混用。
- **Auto 只使用官方模型**,不会自动使用你配置的本地模型。账号没有官方额度时,请手动选择自己配置的模型。
- 插件、代码库索引等 Cursor 功能可以继续使用。
## 数据流转
```text
Cursor 客户端
│
│ Agent 请求与工具结果
▼
cursor-byok 本地服务
│
│ OpenAI / Anthropic 兼容请求
▼
你配置的模型 API
```
API Key、模型配置和应用设置保存在本机。模型请求仍会发送到你选择的上游服务商,请根据对应服务商的隐私政策和计费规则使用。
## 项目结构
```text
cursor-byok/
├── apps/
│ ├── desktop/ # React、Vite、Tauri 桌面应用
│ └── docs/ # Next.js 与 Fumadocs 中文/英文文档站
├── server/ # 本地服务、Cursor 请求处理和模型转发
├── crates/
│ └── semble-core/ # 本地代码索引与搜索核心库
├── cursor-proto/ # Cursor 协议定义
├── benchmarks/ # 代码搜索和索引基准测试
├── images/ # README 展示图片
├── Cargo.toml # Rust 工作区配置
└── Makefile # 常用开发、检查和构建命令
```
## 本地开发
### 环境要求
- Rust 工具链和 Cargo
- Node.js 与 npm
- Tauri 2 的系统构建依赖
- Docker(仅在构建 Docker 镜像时需要)
### 安装依赖
```bash
npm --prefix apps/desktop install
npm --prefix apps/docs install
```
### 启动开发环境
启动桌面前端:
```bash
make dev-web
```
启动桌面应用:
```bash
make dev-desktop
```
启动文档站:
```bash
make dev-docs
```
文档站默认地址为 <http://localhost:3000>。
### 检查与构建
运行完整检查:
```bash
make check
```
分别构建各部分:
```bash
make build-web # 构建桌面前端
make build-server # 构建 Rust 本地服务
make build-docs # 构建文档站
make build-desktop # 构建 Tauri 桌面安装包
make build-docker # 构建 Docker 镜像
```
文档内容位于 `apps/docs/content/docs` 和 `apps/docs/content/blog`。修改文档侧边栏时同步更新 `apps/docs/content/docs/meta.json`。
## 路线图
项目将继续改进模型兼容性、Agent 工具、本地运行稳定性和自托管体验,并探索支持更多 IDE、聊天和 Agent 工作流。
计划与进展请参阅[发布路线图](https://github.com/leookun/cursor-byok/discussions/32)。
## 社区与反馈
- [中文使用指南](https://docs.leokun.cn/zh/docs)
- [GitHub Issues](https://github.com/leookun/cursor-byok/issues)
- [Telegram 社区](https://t.me/cursor_byok)
- QQ 群:`1095916242`、`1094411438`、`1095918002`、`1094419321`
提交问题时,请附上操作系统、cursor-byok 版本、模型类型、请求协议、已脱敏的服务地址、错误信息和复现步骤。请勿公开 API Key 或其他凭据。
## 参与贡献
欢迎提交 Issue 和 Pull Request。提交代码前请先阅读项目中的开发说明,并运行 `make check` 确认格式、测试和前端构建检查通过。
## 许可证
本项目采用 [MIT License](./LICENSE) 开源。
+4 -6
View File
@@ -1,9 +1,3 @@
![Connect cursor-byok to a wide range of model APIs](./images/en-brand-1.png)
![cursor-byok dashboard](./images/en-home-1.png)
<div align="center">
# cursor-byok
@@ -23,6 +17,10 @@ cursor-byok is a local implementation of Cursor's backend.
</div>
![Connect cursor-byok to a wide range of model APIs](./images/en-brand-1.png)
![cursor-byok dashboard](./images/en-home-1.png)
## About
cursor-byok is an open-source local model gateway for Cursor. It runs a service on your machine that connects Cursor to the model APIs you configure, routes model requests through your own providers, and preserves Cursor Agent capabilities such as tool calling, Skills, and MCP.
+35
View File
@@ -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
+3 -3
View File
@@ -21,7 +21,7 @@ use crate::{
store::Store,
};
use super::{inbox::OrderedInbox, CursorCommand, CursorSessionHandle};
use super::{inbox::OrderedInbox, lifecycle, CursorCommand, CursorSessionHandle};
pub struct CursorActor;
@@ -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;
+2 -2
View File
@@ -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);
+1 -1
View File
@@ -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();
+12 -3
View File
@@ -32,17 +32,26 @@ 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 {
handle.cancel();
// 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(())
}
+4
View File
@@ -313,7 +313,11 @@ 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();
tokio::pin!(changed);
changed.as_mut().enable();
if self.inner.runs.lock().await.contains_key(request_id) {
return CursorRoute::Local;
}
+21 -3
View File
@@ -54,8 +54,10 @@ 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"
| "fetchmcpresource" | "getmcptools" => exec::start(runtime, call, context).await,
"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"
| "generateimage" => interaction::start(runtime, call).await,
@@ -66,7 +68,7 @@ pub(super) async fn start(
}
fn normalize_block_until_ms(call: &ToolCall) -> Result<Option<ToolCall>> {
if normalized(&call.name) != "shell" {
if !is_shell_tool(&call.name) {
return Ok(None);
}
let Some(value) = call.arguments.get("block_until_ms") else {
@@ -133,6 +135,10 @@ pub(super) async fn resume_interaction(
interaction::resume(results, search, fetch, pending, response).await
}
fn is_shell_tool(name: &str) -> bool {
matches!(normalized(name).as_str(), "shell" | "bash")
}
pub(super) fn normalized(name: &str) -> String {
name.chars()
.filter(|character| character.is_ascii_alphanumeric())
@@ -167,6 +173,18 @@ mod tests {
assert_eq!(call.arguments["block_until_ms"].as_i64(), Some(45_000));
}
#[test]
fn bash_accepts_integer_valued_float_timeout() {
let call = tool(
"Bash",
serde_json::json!({"command": "echo ok", "block_until_ms": 45_000.0}),
);
let call = normalize_block_until_ms(&call).unwrap().unwrap();
assert_eq!(call.arguments["block_until_ms"].as_i64(), Some(45_000));
}
#[test]
fn shell_rejects_fractional_timeout() {
let call = tool(
+38 -3
View File
@@ -61,6 +61,13 @@ impl Provider for OpenAiChatProvider {
let recorder = self.recorder.clone();
Box::pin(try_stream! {
let ModelInvocation { call_id, request, .. } = invocation;
tracing::debug!(
model = %request.model.model_id,
call_id = %call_id,
history_len = request.history.len(),
tools_count = request.prompt.tools.len(),
"OpenAI Chat provider stream started"
);
let messages = openai_chat_messages(&request.prompt.instructions, &request.history)?;
let mut body = json!({
"model": request.model.model_id,
@@ -85,7 +92,16 @@ impl Provider for OpenAiChatProvider {
_ = cancellation.cancelled() => return,
response = request => response,
};
let response = response?;
let response = match response {
Ok(r) => {
tracing::debug!(status = r.status().as_u16(), "OpenAI Chat HTTP response received");
r
}
Err(e) => {
tracing::debug!(error = %e, "OpenAI Chat HTTP request failed");
Err(Error::from(e))?
}
};
if let Some(recorder) = &recorder {
recorder.response_headers(response.status().as_u16()).await?;
}
@@ -118,15 +134,34 @@ 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() => {
tracing::debug!(
iteration = loop_iteration,
saw_done_marker,
tool_count = tools.len(),
"OpenAI Chat stream cancelled"
);
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 {
tracing::debug!(
iteration = loop_iteration,
saw_done_marker,
"OpenAI Chat SSE stream ended"
);
break;
};
let event = event.map_err(|error| {
let err_msg = error.to_string();
tracing::debug!(iteration = loop_iteration, error = %error, "OpenAI Chat SSE event failed");
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()) {
+52
View File
@@ -90,23 +90,75 @@ 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();
tracing::debug!(
model = %selected,
provider_type = ?provider_type,
timeout_ms = config.request_timeout.as_millis() as u64,
"provider stream created"
);
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(_) => "TextDelta",
super::ModelEvent::TextEnd => "TextEnd",
super::ModelEvent::ThinkingStart => "ThinkingStart",
super::ModelEvent::ThinkingDelta(_) => "ThinkingDelta",
super::ModelEvent::ThinkingEnd => "ThinkingEnd",
super::ModelEvent::ToolCallStart { .. } => "ToolCallStart",
super::ModelEvent::ToolCallArgumentsDelta { .. } => "ToolCallArgsDelta",
super::ModelEvent::ToolCallEnd { .. } => "ToolCallEnd",
super::ModelEvent::ProviderReplayState(_) => "ReplayState",
super::ModelEvent::Usage(_) => "Usage",
super::ModelEvent::Done(_) => "Done",
};
if gap_ms > 5000 {
tracing::debug!(
gap_ms,
elapsed_ms,
event = event_name,
event_count,
"slow gap detected between provider events"
);
}
recorder.event(&event).await?;
last_event_time = now;
yield event;
}
Err(error) => {
tracing::debug!(
error = %error,
elapsed_ms,
gap_ms,
event_count,
"provider stream error"
);
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() {
tracing::debug!(elapsed_ms, event_count, "provider stream ended after cancellation");
recorder.cancelled().await?;
} else {
let error = Error::Provider("provider stream ended without Done".into());
tracing::warn!(
elapsed_ms,
event_count,
"provider stream ended without Done"
);
recorder.failed(&error).await?;
Err(error)?;
}
+32
View File
@@ -19,6 +19,38 @@ use cursor_server::{
};
use prost::Message;
#[tokio::test]
async fn abort_command_cancels_the_run_and_closes_output() {
let (_directory, store) = fixtures::temp_store().await;
let assets = PromptAssets::load(
std::path::Path::new(env!("CARGO_MANIFEST_DIR"))
.join("prompt/cursor")
.as_path(),
)
.unwrap();
let registry = CursorSessionRegistry::new(
store,
Arc::new(fake_provider::FakeProvider::default()),
PromptCompiler::new(assets),
Default::default(),
);
let handle = registry.get_or_create("abort-request").await.unwrap();
let mut output = handle.subscribe();
handle.command(CursorCommand::Abort).await.unwrap();
let frame = tokio::time::timeout(std::time::Duration::from_secs(1), output.recv())
.await
.unwrap()
.expect("Abort must emit a terminal frame");
let (flags, payload) = connect::decode_frames(&frame).unwrap().pop().unwrap();
assert_eq!(flags, connect::END_STREAM_FLAG);
let payload: serde_json::Value = serde_json::from_slice(&payload).unwrap();
assert_eq!(payload["error"]["code"], "canceled");
assert!(handle.cancellation().is_cancelled());
assert_eq!(output.recv().await, None);
}
#[tokio::test]
async fn provider_failure_keeps_the_initial_checkpoint_then_returns_structured_error() {
let (_directory, store) = fixtures::temp_store().await;